From 048d043db4d6ba0c600e8f2350dbdfc8a245e143 Mon Sep 17 00:00:00 2001 From: Ocnrb Date: Fri, 25 Sep 2026 07:51:24 +0100 Subject: [PATCH 1/4] Frame an oversized ADMIN_STATE as a run of chunks in the sync format, locked by shared vectors Co-authored-by: Claude Opus 5.5 --- docs/ADMIN-chunk-vectors.json | 807 ++++++++++++++++++++++ src/js/syncChunks.js | 134 +++- tests/unit/adminChunks.vectors.test.js | 33 + tests/vectors/gen_admin_chunk_vectors.mjs | 97 +++ tests/vectors/regen.mjs | 3 +- 5 files changed, 1035 insertions(+), 39 deletions(-) create mode 100644 docs/ADMIN-chunk-vectors.json create mode 100644 tests/unit/adminChunks.vectors.test.js create mode 100644 tests/vectors/gen_admin_chunk_vectors.mjs diff --git a/docs/ADMIN-chunk-vectors.json b/docs/ADMIN-chunk-vectors.json new file mode 100644 index 0000000..98d2233 --- /dev/null +++ b/docs/ADMIN-chunk-vectors.json @@ -0,0 +1,807 @@ +{ + "limit": 200, + "split": [ + { + "what": "a snapshot that fits travels as itself, with no framing", + "payload": { + "type": "ADMIN_STATE", + "v": 1, + "rev": 4, + "ts": 1789000005000, + "state": { + "bannedMembers": [], + "hiddenMessageIds": [ + "m-1" + ], + "pins": [], + "absorbedThrough": 0 + } + }, + "messages": [ + { + "type": "ADMIN_STATE", + "v": 1, + "rev": 4, + "ts": 1789000005000, + "state": { + "bannedMembers": [], + "hiddenMessageIds": [ + "m-1" + ], + "pins": [], + "absorbedThrough": 0 + } + } + ] + }, + { + "what": "a big snapshot becomes chunks numbered from zero, then its manifest, rev and ts on every row", + "payload": { + "type": "ADMIN_STATE", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "createdBy": "0x1111111111111111111111111111111111111111", + "state": { + "bannedMembers": [ + { + "address": "0x2222222222222222222222222222222222222222", + "sinceEpoch": 3 + } + ], + "hiddenMessageIds": [ + "m-1", + "m-2" + ], + "pins": [ + { + "targetId": "m-9", + "pinnedAt": 1789000000000, + "snapshot": { + "sender": "0x3333333333333333333333333333333333333333", + "text": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + } + } + ], + "absorbedThrough": 0 + } + }, + "messages": [ + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 0, + "chunkCount": 6, + "data": "{\"type\":\"ADMIN_STATE\",\"v\":1,\"rev\":5,\"ts\":1789000009000,\"createdBy\":\"0x1111111111111111111111111111111111111111\",\"state\":{\"bannedMembers\":[{\"address\":\"0x2222222222222222222222222222222222222222\",\"since" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 1, + "chunkCount": 6, + "data": "Epoch\":3}],\"hiddenMessageIds\":[\"m-1\",\"m-2\"],\"pins\":[{\"targetId\":\"m-9\",\"pinnedAt\":1789000000000,\"snapshot\":{\"sender\":\"0x3333333333333333333333333333333333333333\",\"text\":\"xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 2, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 3, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 4, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 5, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx\"}}],\"absorbedThrough\":0}}" + }, + { + "type": "admin_manifest", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkCount": 6 + } + ] + }, + { + "what": "a cut never falls between the two halves of a surrogate pair", + "payload": { + "type": "ADMIN_STATE", + "v": 1, + "rev": 7, + "ts": 1789000029000, + "createdBy": "0x1111111111111111111111111111111111111111", + "state": { + "bannedMembers": [ + { + "address": "0x2222222222222222222222222222222222222222", + "sinceEpoch": 3 + } + ], + "hiddenMessageIds": [ + "m-1", + "m-2" + ], + "pins": [ + { + "targetId": "m-9", + "pinnedAt": 1789000000000, + "snapshot": { + "sender": "0x3333333333333333333333333333333333333333", + "text": "pppppppppppppppppppppppppppppp🐦qqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqq" + } + } + ], + "absorbedThrough": 0 + } + }, + "messages": [ + { + "type": "admin_chunk", + "v": 1, + "rev": 7, + "ts": 1789000029000, + "runId": "runA", + "chunkIndex": 0, + "chunkCount": 4, + "data": "{\"type\":\"ADMIN_STATE\",\"v\":1,\"rev\":7,\"ts\":1789000029000,\"createdBy\":\"0x1111111111111111111111111111111111111111\",\"state\":{\"bannedMembers\":[{\"address\":\"0x2222222222222222222222222222222222222222\",\"since" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 7, + "ts": 1789000029000, + "runId": "runA", + "chunkIndex": 1, + "chunkCount": 4, + "data": "Epoch\":3}],\"hiddenMessageIds\":[\"m-1\",\"m-2\"],\"pins\":[{\"targetId\":\"m-9\",\"pinnedAt\":1789000000000,\"snapshot\":{\"sender\":\"0x3333333333333333333333333333333333333333\",\"text\":\"pppppppppppppppppppppppppppppp" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 7, + "ts": 1789000029000, + "runId": "runA", + "chunkIndex": 2, + "chunkCount": 4, + "data": "🐦qqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqq" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 7, + "ts": 1789000029000, + "runId": "runA", + "chunkIndex": 3, + "chunkCount": 4, + "data": "qqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqq\"}}],\"absorbedThrough\":0}}" + }, + { + "type": "admin_manifest", + "v": 1, + "rev": 7, + "ts": 1789000029000, + "runId": "runA", + "chunkCount": 4 + } + ] + } + ], + "reassemble": [ + { + "what": "a complete run rebuilds the snapshot byte for byte", + "messages": [ + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 0, + "chunkCount": 6, + "data": "{\"type\":\"ADMIN_STATE\",\"v\":1,\"rev\":5,\"ts\":1789000009000,\"createdBy\":\"0x1111111111111111111111111111111111111111\",\"state\":{\"bannedMembers\":[{\"address\":\"0x2222222222222222222222222222222222222222\",\"since" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 1, + "chunkCount": 6, + "data": "Epoch\":3}],\"hiddenMessageIds\":[\"m-1\",\"m-2\"],\"pins\":[{\"targetId\":\"m-9\",\"pinnedAt\":1789000000000,\"snapshot\":{\"sender\":\"0x3333333333333333333333333333333333333333\",\"text\":\"xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 2, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 3, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 4, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 5, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx\"}}],\"absorbedThrough\":0}}" + }, + { + "type": "admin_manifest", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkCount": 6 + } + ], + "payloads": [ + { + "type": "ADMIN_STATE", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "createdBy": "0x1111111111111111111111111111111111111111", + "state": { + "bannedMembers": [ + { + "address": "0x2222222222222222222222222222222222222222", + "sinceEpoch": 3 + } + ], + "hiddenMessageIds": [ + "m-1", + "m-2" + ], + "pins": [ + { + "targetId": "m-9", + "pinnedAt": 1789000000000, + "snapshot": { + "sender": "0x3333333333333333333333333333333333333333", + "text": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + } + } + ], + "absorbedThrough": 0 + } + } + ] + }, + { + "what": "arrival order does not matter", + "messages": [ + { + "type": "admin_manifest", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkCount": 6 + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 5, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx\"}}],\"absorbedThrough\":0}}" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 4, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 3, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 2, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 1, + "chunkCount": 6, + "data": "Epoch\":3}],\"hiddenMessageIds\":[\"m-1\",\"m-2\"],\"pins\":[{\"targetId\":\"m-9\",\"pinnedAt\":1789000000000,\"snapshot\":{\"sender\":\"0x3333333333333333333333333333333333333333\",\"text\":\"xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 0, + "chunkCount": 6, + "data": "{\"type\":\"ADMIN_STATE\",\"v\":1,\"rev\":5,\"ts\":1789000009000,\"createdBy\":\"0x1111111111111111111111111111111111111111\",\"state\":{\"bannedMembers\":[{\"address\":\"0x2222222222222222222222222222222222222222\",\"since" + } + ], + "payloads": [ + { + "type": "ADMIN_STATE", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "createdBy": "0x1111111111111111111111111111111111111111", + "state": { + "bannedMembers": [ + { + "address": "0x2222222222222222222222222222222222222222", + "sinceEpoch": 3 + } + ], + "hiddenMessageIds": [ + "m-1", + "m-2" + ], + "pins": [ + { + "targetId": "m-9", + "pinnedAt": 1789000000000, + "snapshot": { + "sender": "0x3333333333333333333333333333333333333333", + "text": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + } + } + ], + "absorbedThrough": 0 + } + } + ] + }, + { + "what": "two runs in one window stay apart, a whole snapshot among them is not a run", + "messages": [ + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 0, + "chunkCount": 6, + "data": "{\"type\":\"ADMIN_STATE\",\"v\":1,\"rev\":5,\"ts\":1789000009000,\"createdBy\":\"0x1111111111111111111111111111111111111111\",\"state\":{\"bannedMembers\":[{\"address\":\"0x2222222222222222222222222222222222222222\",\"since" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 1, + "chunkCount": 6, + "data": "Epoch\":3}],\"hiddenMessageIds\":[\"m-1\",\"m-2\"],\"pins\":[{\"targetId\":\"m-9\",\"pinnedAt\":1789000000000,\"snapshot\":{\"sender\":\"0x3333333333333333333333333333333333333333\",\"text\":\"xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 2, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 3, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 4, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 5, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx\"}}],\"absorbedThrough\":0}}" + }, + { + "type": "admin_manifest", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkCount": 6 + }, + { + "type": "ADMIN_STATE", + "v": 1, + "rev": 4, + "ts": 1789000005000, + "state": { + "bannedMembers": [], + "hiddenMessageIds": [ + "m-1" + ], + "pins": [], + "absorbedThrough": 0 + } + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 6, + "ts": 1789000019000, + "runId": "runB", + "chunkIndex": 0, + "chunkCount": 6, + "data": "{\"type\":\"ADMIN_STATE\",\"v\":1,\"rev\":6,\"ts\":1789000019000,\"createdBy\":\"0x1111111111111111111111111111111111111111\",\"state\":{\"bannedMembers\":[{\"address\":\"0x2222222222222222222222222222222222222222\",\"since" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 6, + "ts": 1789000019000, + "runId": "runB", + "chunkIndex": 1, + "chunkCount": 6, + "data": "Epoch\":3}],\"hiddenMessageIds\":[\"m-1\",\"m-2\"],\"pins\":[{\"targetId\":\"m-9\",\"pinnedAt\":1789000000000,\"snapshot\":{\"sender\":\"0x3333333333333333333333333333333333333333\",\"text\":\"yyyyyyyyyyyyyyyyyyyyyyyyyyyyyyy" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 6, + "ts": 1789000019000, + "runId": "runB", + "chunkIndex": 2, + "chunkCount": 6, + "data": "yyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyy" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 6, + "ts": 1789000019000, + "runId": "runB", + "chunkIndex": 3, + "chunkCount": 6, + "data": "yyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyy" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 6, + "ts": 1789000019000, + "runId": "runB", + "chunkIndex": 4, + "chunkCount": 6, + "data": "yyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyy" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 6, + "ts": 1789000019000, + "runId": "runB", + "chunkIndex": 5, + "chunkCount": 6, + "data": "yyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyy\"}}],\"absorbedThrough\":0}}" + }, + { + "type": "admin_manifest", + "v": 1, + "rev": 6, + "ts": 1789000019000, + "runId": "runB", + "chunkCount": 6 + } + ], + "payloads": [ + { + "type": "ADMIN_STATE", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "createdBy": "0x1111111111111111111111111111111111111111", + "state": { + "bannedMembers": [ + { + "address": "0x2222222222222222222222222222222222222222", + "sinceEpoch": 3 + } + ], + "hiddenMessageIds": [ + "m-1", + "m-2" + ], + "pins": [ + { + "targetId": "m-9", + "pinnedAt": 1789000000000, + "snapshot": { + "sender": "0x3333333333333333333333333333333333333333", + "text": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + } + } + ], + "absorbedThrough": 0 + } + }, + { + "type": "ADMIN_STATE", + "v": 1, + "rev": 6, + "ts": 1789000019000, + "createdBy": "0x1111111111111111111111111111111111111111", + "state": { + "bannedMembers": [ + { + "address": "0x2222222222222222222222222222222222222222", + "sinceEpoch": 3 + } + ], + "hiddenMessageIds": [ + "m-1", + "m-2" + ], + "pins": [ + { + "targetId": "m-9", + "pinnedAt": 1789000000000, + "snapshot": { + "sender": "0x3333333333333333333333333333333333333333", + "text": "yyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyy" + } + } + ], + "absorbedThrough": 0 + } + } + ] + }, + { + "what": "a run whose head fell out of the window is dropped whole", + "messages": [ + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 1, + "chunkCount": 6, + "data": "Epoch\":3}],\"hiddenMessageIds\":[\"m-1\",\"m-2\"],\"pins\":[{\"targetId\":\"m-9\",\"pinnedAt\":1789000000000,\"snapshot\":{\"sender\":\"0x3333333333333333333333333333333333333333\",\"text\":\"xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 2, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 3, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 4, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 5, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx\"}}],\"absorbedThrough\":0}}" + }, + { + "type": "admin_manifest", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkCount": 6 + } + ], + "payloads": [] + }, + { + "what": "a run with no manifest is not assumed complete", + "messages": [ + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 0, + "chunkCount": 6, + "data": "{\"type\":\"ADMIN_STATE\",\"v\":1,\"rev\":5,\"ts\":1789000009000,\"createdBy\":\"0x1111111111111111111111111111111111111111\",\"state\":{\"bannedMembers\":[{\"address\":\"0x2222222222222222222222222222222222222222\",\"since" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 1, + "chunkCount": 6, + "data": "Epoch\":3}],\"hiddenMessageIds\":[\"m-1\",\"m-2\"],\"pins\":[{\"targetId\":\"m-9\",\"pinnedAt\":1789000000000,\"snapshot\":{\"sender\":\"0x3333333333333333333333333333333333333333\",\"text\":\"xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 2, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 3, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 4, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx" + }, + { + "type": "admin_chunk", + "v": 1, + "rev": 5, + "ts": 1789000009000, + "runId": "runA", + "chunkIndex": 5, + "chunkCount": 6, + "data": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx\"}}],\"absorbedThrough\":0}}" + } + ], + "payloads": [] + }, + { + "what": "a run whose chunks do not form JSON is dropped", + "messages": [ + { + "type": "admin_chunk", + "v": 1, + "rev": 1, + "ts": 1, + "runId": "bad", + "chunkIndex": 0, + "chunkCount": 1, + "data": "{oops" + }, + { + "type": "admin_manifest", + "v": 1, + "rev": 1, + "ts": 1, + "runId": "bad", + "chunkCount": 1 + } + ], + "payloads": [] + } + ] +} diff --git a/src/js/syncChunks.js b/src/js/syncChunks.js index 49f9f09..046516d 100644 --- a/src/js/syncChunks.js +++ b/src/js/syncChunks.js @@ -3,7 +3,8 @@ * chunks closed by a manifest, the same shape image blobs use. The network * drops an oversized message WITHOUT an error the publisher can see. * - * Android parity: core/SyncChunks.kt, locked by tests/vectors/sync-chunks.json. + * Android parity: core/SyncChunks.kt, locked by docs/SYNC-chunk-vectors.json + * and docs/ADMIN-chunk-vectors.json. */ /** @@ -13,80 +14,137 @@ */ export const SYNC_CHUNK_CHARS = 150 * 1024; +/** + * The row types and fields of one framed protocol. `carry` names the payload + * fields every row repeats; `keepPairs` never cuts between the two halves of + * a surrogate pair, which a UTF-8 encoder downstream would turn into '?'. + */ +const SYNC_FRAME = Object.freeze({ + chunk: 'sync_chunk', manifest: 'sync_manifest', id: 'syncId', carry: ['ts'], keepPairs: false +}); + +/** + * An ADMIN_STATE too big for one message. Its own row types, so a reader that + * predates the split skips the rows instead of misreading them; `rev` rides + * every row so a run can be ranked before it is assembled. + */ +export const ADMIN_FRAME = Object.freeze({ + chunk: 'admin_chunk', manifest: 'admin_manifest', id: 'runId', carry: ['rev', 'ts'], keepPairs: true +}); + /** * Frame a payload for the wire. * - * @param {Object} payload - The whole `{ type:'sync', v, ts, data }` snapshot - * @param {string} syncId - Run id, tying the chunks to their manifest - * @param {number} [limit] - Characters per chunk + * @param {Object} payload - The whole snapshot + * @param {string} runId - Ties the chunks to their manifest + * @param {Object} frame - Row types and fields (SYNC_FRAME, ADMIN_FRAME) + * @param {number} limit - Characters per chunk * @returns {Object[]} - The payload itself when it fits, else chunks + manifest */ -export function splitSyncPayload(payload, syncId, limit = SYNC_CHUNK_CHARS) { +export function splitFramed(payload, runId, frame, limit) { const serialised = JSON.stringify(payload); if (serialised.length <= limit) return [payload]; - const chunkCount = Math.ceil(serialised.length / limit); - const out = []; - for (let i = 0; i < chunkCount; i++) { - out.push({ - type: 'sync_chunk', v: 1, ts: payload.ts, syncId, - chunkIndex: i, chunkCount, - data: serialised.slice(i * limit, (i + 1) * limit) - }); + const slices = []; + for (let start = 0; start < serialised.length;) { + let end = Math.min(start + limit, serialised.length); + if (frame.keepPairs && end < serialised.length && end - 1 > start + && isHighSurrogate(serialised.charCodeAt(end - 1))) { + end -= 1; + } + slices.push(serialised.slice(start, end)); + start = end; } + + const header = {}; + for (const field of frame.carry) header[field] = payload[field]; + const chunkCount = slices.length; + const out = slices.map((data, chunkIndex) => ({ + type: frame.chunk, v: 1, ...header, [frame.id]: runId, chunkIndex, chunkCount, data + })); // Last, so a reader holding the manifest knows the run is complete rather // than still arriving. - out.push({ type: 'sync_manifest', v: 1, ts: payload.ts, syncId, chunkCount }); + out.push({ type: frame.manifest, v: 1, ...header, [frame.id]: runId, chunkCount }); return out; } /** - * Put the wire back together. An incomplete run is dropped whole, never - * applied in part: a truncated snapshot would merge garbage into the account. + * Put the runs back together. An incomplete run is dropped whole, never + * applied in part: a truncated snapshot would merge garbage into the state. + * Rows of any other type are ignored; what the joined payload must be is the + * caller's to check. * * @param {Object[]} messages - Opened payloads, in any order + * @param {Object} frame - Row types and fields * @param {(reason: Object) => void} [onDropped] - Told about each dropped run - * @returns {Object[]} - Whole `sync` payloads + * @returns {Array<{runId: string, manifest: Object, payload: Object}>} */ -export function reassembleSyncPayloads(messages, onDropped) { - const whole = []; - const parts = new Map(); // syncId -> Map(index, slice) - const expected = new Map(); // syncId -> chunkCount +export function joinFramed(messages, frame, onDropped) { + const parts = new Map(); // runId -> Map(index, slice) + const manifests = new Map(); // runId -> manifest row for (const m of messages || []) { if (!m || m.v !== 1) continue; - if (m.type === 'sync') { - whole.push(m); - } else if (m.type === 'sync_chunk' - && typeof m.syncId === 'string' && m.syncId + const runId = m[frame.id]; + if (typeof runId !== 'string' || !runId) continue; + if (m.type === frame.chunk && typeof m.data === 'string' && Number.isInteger(m.chunkIndex) && m.chunkIndex >= 0) { - if (!parts.has(m.syncId)) parts.set(m.syncId, new Map()); - parts.get(m.syncId).set(m.chunkIndex, m.data); - } else if (m.type === 'sync_manifest' - && typeof m.syncId === 'string' && m.syncId) { - expected.set(m.syncId, m.chunkCount); + if (!parts.has(runId)) parts.set(runId, new Map()); + parts.get(runId).set(m.chunkIndex, m.data); + } else if (m.type === frame.manifest) { + manifests.set(runId, m); } } - for (const [syncId, chunkCount] of expected) { - const got = parts.get(syncId); + const out = []; + for (const [runId, manifest] of manifests) { + const chunkCount = manifest.chunkCount; + const got = parts.get(runId); if (!Number.isInteger(chunkCount) || chunkCount <= 0 || !got || got.size !== chunkCount) { - onDropped?.({ syncId, have: got?.size ?? 0, want: chunkCount, reason: 'incomplete' }); + onDropped?.({ [frame.id]: runId, have: got?.size ?? 0, want: chunkCount, reason: 'incomplete' }); continue; } let joined = ''; for (let i = 0; i < chunkCount; i++) joined += got.get(i) ?? ''; - let payload; try { - payload = JSON.parse(joined); + out.push({ runId, manifest, payload: JSON.parse(joined) }); } catch (e) { - onDropped?.({ syncId, reason: 'unparseable', error: e.message }); - continue; + onDropped?.({ [frame.id]: runId, reason: 'unparseable', error: e.message }); } + } + return out; +} + +function isHighSurrogate(code) { + return code >= 0xD800 && code <= 0xDBFF; +} + +/** + * Frame a sync snapshot for the wire. + * + * @param {Object} payload - The whole `{ type:'sync', v, ts, data }` snapshot + * @param {string} syncId - Run id, tying the chunks to their manifest + * @param {number} [limit] - Characters per chunk + * @returns {Object[]} - The payload itself when it fits, else chunks + manifest + */ +export function splitSyncPayload(payload, syncId, limit = SYNC_CHUNK_CHARS) { + return splitFramed(payload, syncId, SYNC_FRAME, limit); +} + +/** + * Put the sync wire back together. + * + * @param {Object[]} messages - Opened payloads, in any order + * @param {(reason: Object) => void} [onDropped] - Told about each dropped run + * @returns {Object[]} - Whole `sync` payloads + */ +export function reassembleSyncPayloads(messages, onDropped) { + const whole = (messages || []).filter(m => m && m.v === 1 && m.type === 'sync'); + for (const { runId, payload } of joinFramed(messages, SYNC_FRAME, onDropped)) { if (payload?.type === 'sync' && payload.v === 1) whole.push(payload); - else onDropped?.({ syncId, reason: 'not a snapshot' }); + else onDropped?.({ syncId: runId, reason: 'not a snapshot' }); } return whole; } diff --git a/tests/unit/adminChunks.vectors.test.js b/tests/unit/adminChunks.vectors.test.js new file mode 100644 index 0000000..fac2571 --- /dev/null +++ b/tests/unit/adminChunks.vectors.test.js @@ -0,0 +1,33 @@ +// The ADMIN_STATE framing vectors are the shared spec: this suite and +// AdminChunksTest.kt read the same JSON, so a snapshot split by one client's +// owner reassembles on the other client's members. +import { describe, it, expect } from 'vitest'; +import { readFileSync } from 'fs'; +import { fileURLToPath } from 'url'; +import { dirname, join } from 'path'; +import { splitFramed, joinFramed, ADMIN_FRAME } from '../../src/js/syncChunks.js'; + +const vectors = JSON.parse(readFileSync( + join(dirname(fileURLToPath(import.meta.url)), '..', '..', + 'docs', 'ADMIN-chunk-vectors.json'), 'utf8')); + +describe('ADMIN_STATE framing parity vectors', () => { + for (const v of vectors.split) { + it(v.what, () => { + expect(splitFramed(v.payload, 'runA', ADMIN_FRAME, vectors.limit)).toEqual(v.messages); + }); + } + + for (const v of vectors.reassemble) { + it(v.what, () => { + expect(joinFramed(v.messages, ADMIN_FRAME).map(r => r.payload)).toEqual(v.payloads); + }); + } + + it('keeps each chunk free of a lone surrogate half', () => { + const straddling = vectors.split.find(v => v.what.includes('surrogate')); + for (const m of straddling.messages.filter(r => r.type === 'admin_chunk')) { + expect(m.data).toBe(m.data.toWellFormed()); + } + }); +}); diff --git a/tests/vectors/gen_admin_chunk_vectors.mjs b/tests/vectors/gen_admin_chunk_vectors.mjs new file mode 100644 index 0000000..c045908 --- /dev/null +++ b/tests/vectors/gen_admin_chunk_vectors.mjs @@ -0,0 +1,97 @@ +// One-shot generator for the ADMIN_STATE framing parity vectors. +// +// A moderation snapshot too big for one wire message travels on the -3 as a +// run of admin_chunk rows closed by an admin_manifest. Both clients must frame +// it the same way, or a snapshot split by one owner device never reassembles +// on the members' other client, and they keep the previous moderation. +// +// The vectors fix the split (how a snapshot becomes rows, at a small budget +// so the fixtures stay readable: rev and ts on every row, a surrogate pair +// never cut) and the reassembly rules (order does not matter, runs are kept +// apart, an incomplete or unparseable run is dropped whole). +import { splitFramed, joinFramed, ADMIN_FRAME } from '../../src/js/syncChunks.js'; + +const LIMIT = 200; + +const snapshot = (rev, ts, pinText) => ({ + type: 'ADMIN_STATE', v: 1, rev, ts, createdBy: '0x1111111111111111111111111111111111111111', + state: { + bannedMembers: [{ address: '0x2222222222222222222222222222222222222222', sinceEpoch: 3 }], + hiddenMessageIds: ['m-1', 'm-2'], + pins: [{ targetId: 'm-9', pinnedAt: 1789000000000, snapshot: { sender: '0x3333333333333333333333333333333333333333', text: pinText } }], + absorbedThrough: 0 + } +}); + +const small = { + type: 'ADMIN_STATE', v: 1, rev: 4, ts: 1789000005000, + state: { bannedMembers: [], hiddenMessageIds: ['m-1'], pins: [], absorbedThrough: 0 } +}; +const big = snapshot(5, 1789000009000, 'x'.repeat(700)); +const other = snapshot(6, 1789000019000, 'y'.repeat(700)); + +// An emoji placed so a cut at a multiple of the budget would fall between its +// two halves. +const prefix = JSON.stringify(snapshot(7, 1789000029000, '')).indexOf('"text":""') + '"text":"'.length; +const cut = Math.ceil((prefix + 1) / LIMIT) * LIMIT; +const straddling = snapshot(7, 1789000029000, 'p'.repeat(cut - 1 - prefix) + '\u{1F426}' + 'q'.repeat(300)); + +const bigRun = splitFramed(big, 'runA', ADMIN_FRAME, LIMIT); +const otherRun = splitFramed(other, 'runB', ADMIN_FRAME, LIMIT); +const payloadsOf = (messages) => joinFramed(messages, ADMIN_FRAME).map(r => r.payload); + +console.log(JSON.stringify({ + limit: LIMIT, + split: [ + { + what: 'a snapshot that fits travels as itself, with no framing', + payload: small, + messages: splitFramed(small, 'runA', ADMIN_FRAME, LIMIT) + }, + { + what: 'a big snapshot becomes chunks numbered from zero, then its manifest, rev and ts on every row', + payload: big, + messages: bigRun + }, + { + what: 'a cut never falls between the two halves of a surrogate pair', + payload: straddling, + messages: splitFramed(straddling, 'runA', ADMIN_FRAME, LIMIT) + } + ], + reassemble: [ + { + what: 'a complete run rebuilds the snapshot byte for byte', + messages: bigRun, + payloads: [big] + }, + { + what: 'arrival order does not matter', + messages: [...bigRun].reverse(), + payloads: [big] + }, + { + what: 'two runs in one window stay apart, a whole snapshot among them is not a run', + messages: [...bigRun, small, ...otherRun], + payloads: payloadsOf([...bigRun, small, ...otherRun]) + }, + { + what: 'a run whose head fell out of the window is dropped whole', + messages: bigRun.slice(1), + payloads: [] + }, + { + what: 'a run with no manifest is not assumed complete', + messages: bigRun.filter(m => m.type !== 'admin_manifest'), + payloads: [] + }, + { + what: 'a run whose chunks do not form JSON is dropped', + messages: [ + { type: 'admin_chunk', v: 1, rev: 1, ts: 1, runId: 'bad', chunkIndex: 0, chunkCount: 1, data: '{oops' }, + { type: 'admin_manifest', v: 1, rev: 1, ts: 1, runId: 'bad', chunkCount: 1 } + ], + payloads: [] + } + ] +}, null, 2)); diff --git a/tests/vectors/regen.mjs b/tests/vectors/regen.mjs index 0e723b5..79bf204 100644 --- a/tests/vectors/regen.mjs +++ b/tests/vectors/regen.mjs @@ -34,7 +34,8 @@ const GENERATORS = { 'gen_storage_read_vectors.mjs': 'STORAGE-signed-read-vectors.json', 'gen_storage_purge_vectors.mjs': 'STORAGE-purge-vectors.json', 'gen_storage_stored_vectors.mjs': 'STORAGE-stored-vectors.json', - 'gen_sync_merge_vectors.mjs': 'SYNC-merge-vectors.json' + 'gen_sync_merge_vectors.mjs': 'SYNC-merge-vectors.json', + 'gen_admin_chunk_vectors.mjs': 'ADMIN-chunk-vectors.json' }; const androidDocs = (() => { From 3167f7810515ad8f35e3f2eb44319e37b6db3eef Mon Sep 17 00:00:00 2001 From: Ocnrb Date: Fri, 25 Sep 2026 07:51:24 +0100 Subject: [PATCH 2/4] Split an ADMIN_STATE that does not fit one wire message and refuse it past four chunks Co-authored-by: Claude Opus 5.5 --- src/js/channels/AdminState.js | 102 +++++-- src/js/channels/AdminStateConfirm.js | 6 + src/js/channels/MessageFlow.js | 4 + src/js/channels/StorageCopy.js | 4 +- src/js/config.js | 8 + src/js/crypto.js | 11 + src/js/streamr.js | 256 ++++++++++++------ src/js/subscriptionManager.js | 4 + src/js/ui/ChannelSettingsUI.js | 2 + tests/unit/adminState.split.test.js | 172 ++++++++++++ tests/unit/adminStateConfirm.test.js | 11 + .../unit/messageFlow.adminInvalidate.test.js | 62 +++++ tests/unit/storageCopy.test.js | 26 +- tests/unit/streamr.adminRuns.test.js | 148 ++++++++++ tests/unit/streamr.wireBytes.test.js | 134 +++++++++ 15 files changed, 842 insertions(+), 108 deletions(-) create mode 100644 tests/unit/adminState.split.test.js create mode 100644 tests/unit/messageFlow.adminInvalidate.test.js create mode 100644 tests/unit/streamr.adminRuns.test.js create mode 100644 tests/unit/streamr.wireBytes.test.js diff --git a/src/js/channels/AdminState.js b/src/js/channels/AdminState.js index 116830c..9257962 100644 --- a/src/js/channels/AdminState.js +++ b/src/js/channels/AdminState.js @@ -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 { /** @@ -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 } /** @@ -382,7 +386,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); @@ -409,32 +422,73 @@ 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. Signals arriving while a + * read is pending share it. + */ + readAfterSignal(messageStreamId) { + if (this._signalReads.has(messageStreamId)) return; + this._signalReads.set(messageStreamId, setTimeout(() => { + this._signalReads.delete(messageStreamId); + if (adminStatePoller.getStreamId() === messageStreamId) adminStatePoller.pollNow(); + }, CONFIG.subscriptions.adminSignalReadDelayMs)); } // High-level convenience helpers built on top of publishAdminState ---------- diff --git a/src/js/channels/AdminStateConfirm.js b/src/js/channels/AdminStateConfirm.js index 6dd30be..d6e1d97 100644 --- a/src/js/channels/AdminStateConfirm.js +++ b/src/js/channels/AdminStateConfirm.js @@ -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; } diff --git a/src/js/channels/MessageFlow.js b/src/js/channels/MessageFlow.js index 8cf0e20..a58a305 100644 --- a/src/js/channels/MessageFlow.js +++ b/src/js/channels/MessageFlow.js @@ -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); + return; + } if (!this.manager._isValidAdminState(data.snapshot)) return; // Sender authenticity is already validated above (account === diff --git a/src/js/channels/StorageCopy.js b/src/js/channels/StorageCopy.js index b643ed7..2620c12 100644 --- a/src/js/channels/StorageCopy.js +++ b/src/js/channels/StorageCopy.js @@ -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)) { diff --git a/src/js/config.js b/src/js/config.js index e6521ee..cb067f9 100644 --- a/src/js/config.js +++ b/src/js/config.js @@ -293,6 +293,14 @@ 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 this late, + // once storage has had time to hold the run. + adminSignalReadDelayMs: 5000, // 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], diff --git a/src/js/crypto.js b/src/js/crypto.js index 618622a..26c8e88 100644 --- a/src/js/crypto.js +++ b/src/js/crypto.js @@ -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 diff --git a/src/js/streamr.js b/src/js/streamr.js index 12a7712..7f76e6d 100644 --- a/src/js/streamr.js +++ b/src/js/streamr.js @@ -61,6 +61,7 @@ import { STREAM_CONFIG } from './streamConfig.js'; import { History } from './streamr/History.js'; import { MessagePipeline } from './streamr/MessagePipeline.js'; import { storageFetch } from './storageFetch.js'; +import { ADMIN_FRAME, joinFramed } from './syncChunks.js'; // === ID DERIVATION FUNCTIONS === // Re-exported from streamConstants.js; kept as local names for readability @@ -2743,6 +2744,59 @@ class StreamrController { return new TextEncoder().encode(JSON.stringify(envelope)).length; } + /** + * Bytes [publishAdminState] would hand the transport for this snapshot, + * branch for branch, including the publish's one attempt to recover a + * missing epoch key. A test holds the two together. + */ + async adminWireBytes(adminStreamId, state, password = null) { + try { + const { channelManager } = await import('./channels.js'); + const base = String(adminStreamId).replace(/-[12345]$/, ''); + const channel = channelManager.channels?.get(base + '-1'); + if (channel?.type === 'gated' || channel?.gate?.address) { + const { epochKeyManager } = await import('./epochKeyManager.js'); + if (!await epochKeyManager.getCurrentKey(channel.messageStreamId)) { + await epochKeyManager.ensureChannelKeys(channel); + } + return await this.epochWireBytes(channel, adminStreamId, state); + } + } catch (e) { + if (String(e?.message || '').includes('No epoch key')) throw e; + /* registry unavailable → legacy path */ + } + return this._plainWireBytes(state, password); + } + + /** + * Bytes [publishAsChannel] would hand the transport for this payload, + * branch for branch. A test holds the two together. + */ + async channelWireBytes(streamId, data, password = null) { + let channelManager = null; + try { + ({ channelManager } = await import('./channels.js')); + } catch { /* registry unavailable → ephemeral (public/password) */ } + const base = String(streamId).replace(/-[12345]$/, ''); + const record = channelManager?.channels?.get(base + '-1') + ?? (channelManager?.previewChannel?.messageStreamId === base + '-1' + ? channelManager.previewChannel : null); + if (record?.type === 'gated' || record?.gate?.address) { + return this.epochWireBytes(record, streamId, data); + } + if (channelManager?.usesAccountPublish?.(streamId)) { + return this._plainWireBytes(stripLocalFields(data), password); + } + const { proof } = getChannelIdentity(streamId); + return this._plainWireBytes({ ...stripLocalFields(data), proof }, password); + } + + /** JSON as the SDK serialises it; with a password, the ciphertext string that replaces it. */ + _plainWireBytes(payload, password) { + const bytes = new TextEncoder().encode(JSON.stringify(payload)).length; + return password ? cryptoManager.encryptedLength(bytes) + 2 : bytes; + } + /** * Re-key a Sealed channel's shared publish grants: the new key's * address gains publish+subscribe on -1/-2 and the old one loses @@ -3483,98 +3537,142 @@ class StreamrController { } const last = historyCount ?? STREAM_CONFIG.ADMIN_HISTORY_COUNT; - const partition = STREAM_CONFIG.ADMIN_STREAM.MODERATION; - - let latest = null; + let read; try { - // Gated: raw resend — the SDK validator re-checks stored envelopes - // against the present gate state; authorship is established - // client-side by resolveAuthor (admin-only on -3) instead. - const gatedChannel = await this._gatedChannelFor(adminStreamId); - const resend = await this.client.resend( - { streamId: adminStreamId, partition }, - { last, raw: true } - ); + read = await this._readAdminWindow(adminStreamId, last, password); + // The newest snapshot went out split and its run reaches further + // back than this window: read once more, wide enough to hold it. + if (read.cutRun) { + read = await this._readAdminWindow( + adminStreamId, last + read.cutRun.chunkCount + 1, password); + } + } catch (error) { + Logger.warn('resendAdminState error:', error.message); + return null; + } - const iterator = resend[Symbol.asyncIterator](); - let iteratorDone = false; + Logger.debug('resendAdminState result:', { + adminStreamId: adminStreamId.slice(-30), + found: !!read.latest, + rev: read.latest?.rev + }); + return read.latest; + } - while (!iteratorDone) { - let message; - try { - const result = await iterator.next(); - iteratorDone = result.done; - if (iteratorDone) break; - message = result.value; - } catch (iterError) { - if (iterError.code === 'DECRYPT_ERROR' || iterError.message?.includes('encryption key')) { - continue; - } - Logger.warn('resendAdminState iteration error:', iterError.message); - continue; - } + /** + * One resend of the -3/P0 window: the newest snapshot among whole + * ADMIN_STATE rows and complete runs, plus the manifest of a newer run + * the window cut short, if any. Rows only join a run of their own + * publisher, and each row passes the same authority check a whole + * snapshot does. + * @private + */ + async _readAdminWindow(adminStreamId, last, password) { + const partition = STREAM_CONFIG.ADMIN_STREAM.MODERATION; + let latest = null; + const framed = new Map(); // publisherId -> opened chunk/manifest rows + const num = (v) => (typeof v === 'number' ? v : 0); + const newer = (a, b) => !b + || num(a.rev) > num(b.rev) + || (num(a.rev) === num(b.rev) && num(a.ts) > num(b.ts)); + + // Gated: raw resend — the SDK validator re-checks stored envelopes + // against the present gate state; authorship is established + // client-side by resolveAuthor (admin-only on -3) instead. + const gatedChannel = await this._gatedChannelFor(adminStreamId); + const resend = await this.client.resend( + { streamId: adminStreamId, partition }, + { last, raw: true } + ); - try { - if (!gatedChannel && (!verifyEnvelopeAuthenticity(message) - || !await this.publisherMayWrite(adminStreamId, message))) continue; - let content = message.content || message; - if (password && typeof content === 'string') { - try { - content = await cryptoManager.decryptJSON(content, password); - } catch (decryptError) { - continue; - } - } - // Gated: ADMIN_STATE arrives as an epoch envelope. History - // context so entries sealed under an older epoch open in - // that epoch's validity window instead of being dropped. - if (this.isEpochEnvelope(content)) { - const judged = storageFetch.judgeMessage(adminStreamId, partition, message); - if (judged.forwardDated) continue; - const opened = await this.openEpochEnvelope(adminStreamId, content, - { live: false, timestamp: judged.judgeTime }); - if (opened === null) continue; - content = opened; - } - if (!content || typeof content !== 'object') continue; - if (content.type && content.type !== 'ADMIN_STATE') continue; + const iterator = resend[Symbol.asyncIterator](); + let iteratorDone = false; - // Inject publisher info so caller can validate sender == createdBy. - // Gated: the clone publishes for everyone — resolveAuthor swaps in - // the envelope signer and DROPS non-admin writes on -3 (D10c). - { - const transportPublisher = typeof message.getPublisherId === 'function' - ? message.getPublisherId() : message.publisherId; - const publisherId = await this.resolveAuthor( - adminStreamId, message, transportPublisher); - if (!publisherId) continue; - if (!content.createdBy) content.createdBy = publisherId; - } + while (!iteratorDone) { + let message; + try { + const result = await iterator.next(); + iteratorDone = result.done; + if (iteratorDone) break; + message = result.value; + } catch (iterError) { + if (iterError.code === 'DECRYPT_ERROR' || iterError.message?.includes('encryption key')) { + continue; + } + Logger.warn('resendAdminState iteration error:', iterError.message); + continue; + } - const incomingRev = typeof content.rev === 'number' ? content.rev : 0; - const incomingTs = typeof content.ts === 'number' ? content.ts : 0; - const latestRev = latest ? (latest.rev || 0) : -1; - const latestTs = latest ? (latest.ts || 0) : 0; - if (incomingRev > latestRev || (incomingRev === latestRev && incomingTs > latestTs)) { - latest = content; + try { + if (!gatedChannel && (!verifyEnvelopeAuthenticity(message) + || !await this.publisherMayWrite(adminStreamId, message))) continue; + let content = message.content || message; + if (password && typeof content === 'string') { + try { + content = await cryptoManager.decryptJSON(content, password); + } catch (decryptError) { + continue; } - } catch (e) { - Logger.debug('resendAdminState entry processing error:', e.message); + } + // Gated: ADMIN_STATE arrives as an epoch envelope. History + // context so entries sealed under an older epoch open in + // that epoch's validity window instead of being dropped. + if (this.isEpochEnvelope(content)) { + const judged = storageFetch.judgeMessage(adminStreamId, partition, message); + if (judged.forwardDated) continue; + const opened = await this.openEpochEnvelope(adminStreamId, content, + { live: false, timestamp: judged.judgeTime }); + if (opened === null) continue; + content = opened; + } + if (!content || typeof content !== 'object') continue; + const isFramed = content.type === ADMIN_FRAME.chunk || content.type === ADMIN_FRAME.manifest; + if (content.type && content.type !== 'ADMIN_STATE' && !isFramed) continue; + + // Inject publisher info so caller can validate sender == createdBy. + // Gated: the clone publishes for everyone — resolveAuthor swaps in + // the envelope signer and DROPS non-admin writes on -3 (D10c). + const transportPublisher = typeof message.getPublisherId === 'function' + ? message.getPublisherId() : message.publisherId; + const publisherId = await this.resolveAuthor( + adminStreamId, message, transportPublisher); + if (!publisherId) continue; + + if (isFramed) { + if (!framed.has(publisherId)) framed.set(publisherId, []); + framed.get(publisherId).push(content); continue; } + if (!content.createdBy) content.createdBy = publisherId; + if (newer(content, latest)) latest = content; + } catch (e) { + Logger.debug('resendAdminState entry processing error:', e.message); + continue; } - } catch (error) { - Logger.warn('resendAdminState error:', error.message); - return null; } - Logger.debug('resendAdminState result:', { - adminStreamId: adminStreamId.slice(-30), - found: !!latest, - rev: latest?.rev - }); - return latest; + let cutRun = null; + for (const [publisherId, rows] of framed) { + const dropped = []; + for (const { manifest, payload } of joinFramed(rows, ADMIN_FRAME, (d) => dropped.push(d))) { + if (payload?.type !== 'ADMIN_STATE' + || payload.rev !== manifest.rev || payload.ts !== manifest.ts) continue; + if (!payload.createdBy) payload.createdBy = publisherId; + if (newer(payload, latest)) latest = payload; + } + for (const d of dropped) { + if (d.reason !== 'incomplete') continue; + const manifest = rows.find(r => r.type === ADMIN_FRAME.manifest && r.runId === d.runId); + if (manifest && Number.isInteger(manifest.chunkCount) && newer(manifest, cutRun)) { + cutRun = manifest; + } + } + } + if (cutRun && (!newer(cutRun, latest) + || cutRun.chunkCount + 1 <= last + || cutRun.chunkCount > CONFIG.subscriptions.adminReadMaxChunks)) cutRun = null; + return { latest, cutRun }; } /** diff --git a/src/js/subscriptionManager.js b/src/js/subscriptionManager.js index 987ee81..94c36c0 100644 --- a/src/js/subscriptionManager.js +++ b/src/js/subscriptionManager.js @@ -586,6 +586,10 @@ class SubscriptionManager { // can publish moderation), and the matching -2/P0 signal is // published by the same admin. We additionally inject createdBy // from the account so applyAdminState's owner check passes. + if (msg.snapshot === undefined) { + channelManager.adminState?.readAfterSignal(streamId); + return; + } if (!msg.snapshot || typeof msg.snapshot !== 'object') return; const sender = (msg.account || msg.user || '').toLowerCase(); if (!msg.snapshot.createdBy && sender) { diff --git a/src/js/ui/ChannelSettingsUI.js b/src/js/ui/ChannelSettingsUI.js index d7980ca..b0b621a 100644 --- a/src/js/ui/ChannelSettingsUI.js +++ b/src/js/ui/ChannelSettingsUI.js @@ -693,6 +693,8 @@ class ChannelSettingsUI { const { showNotification } = this.deps; if (event === 'admin_state_unconfirmed') { showNotification?.('Moderation change not yet confirmed on storage. It will be retried when you open the channel again.', 'warning'); + } else if (event === 'admin_state_too_large') { + showNotification?.("This channel's moderation state is too large to publish. Unpin some messages and try again.", 'warning'); } else if (event === 'admin_state_superseded') { showNotification?.('Moderation was changed from another device; the last change made here was replaced.', 'warning'); } else if (event === 'storage_copy_confirmed') { diff --git a/tests/unit/adminState.split.test.js b/tests/unit/adminState.split.test.js new file mode 100644 index 0000000..0e13ac4 --- /dev/null +++ b/tests/unit/adminState.split.test.js @@ -0,0 +1,172 @@ +/** + * Publishing an ADMIN_STATE that no longer fits one wire message. + * + * Past the DataChannel's max-message-size the SDK throws where the publisher + * never sees it and the snapshot is simply gone. So the snapshot is measured + * as it will travel: whole when it fits, as it always was; split into a run + * when it does not; refused, with the owner told, past the cap. + */ + +import { describe, it, expect, beforeEach, afterEach, vi } from 'vitest'; + +vi.mock('../../src/js/logger.js', () => ({ + Logger: { debug: vi.fn(), info: vi.fn(), warn: vi.fn(), error: vi.fn() } +})); + +const OWNER = '0x' + 'aa'.repeat(20); + +vi.mock('../../src/js/auth.js', () => ({ + authManager: { getAddress: () => '0x' + 'aa'.repeat(20) } +})); + +/** Epoch-envelope-like inflation: base64 over the ciphertext. */ +const inflate = (row) => Math.ceil(JSON.stringify(row).length * 1.4); + +vi.mock('../../src/js/streamr.js', () => ({ + streamrController: { + adminWireBytes: vi.fn(async (_id, row) => inflate(row)), + channelWireBytes: vi.fn(async (_id, signal) => inflate(signal)), + publishAdminState: vi.fn(async () => ({ timestamp: Date.now() })), + publishControl: vi.fn(async () => ({ timestamp: Date.now() })) + }, + STREAM_CONFIG: { ADMIN_HISTORY_COUNT: 10 }, + deriveEphemeralId: (id) => id.replace(/-1$/, '-2'), + deriveAdminId: (id) => id.replace(/-1$/, '-3') +})); + +const { AdminState } = await import('../../src/js/channels/AdminState.js'); +const { streamrController } = await import('../../src/js/streamr.js'); +const { adminStatePoller } = await import('../../src/js/adminStatePoller.js'); +const { CONFIG } = await import('../../src/js/config.js'); +const { joinFramed, ADMIN_FRAME } = await import('../../src/js/syncChunks.js'); + +const STREAM = `${OWNER}/room-1`; + +describe('ADMIN_STATE publish on the wire', () => { + let admin; + let manager; + let channel; + + beforeEach(() => { + vi.clearAllMocks(); + channel = { + messageStreamId: STREAM, + createdBy: OWNER, + adminLoaded: true, + adminRev: 7, + adminSnapshot: { bannedMembers: [], hiddenMessageIds: [], pins: [], absorbedThrough: 0 } + }; + manager = { + channels: new Map([[STREAM, channel]]), + notifyHandlers: vi.fn(), + adminConfirm: { track: vi.fn() } + }; + admin = new AdminState(manager); + for (const m of ['_isValidAdminState', '_normalizeAdminState', 'applyAdminState', '_publishAdminStateInner']) { + manager[m] = admin[m].bind(admin); + } + }); + + const pinOf = (chars) => ({ targetId: 'm-1', pinnedAt: 1, snapshot: { text: 'p'.repeat(chars) } }); + const published = () => streamrController.publishAdminState.mock.calls.map(c => c[1]); + const signalled = async () => { + await vi.waitFor(() => expect(streamrController.publishControl).toHaveBeenCalled()); + return streamrController.publishControl.mock.calls[0][1]; + }; + + it('sends a snapshot that fits as one message, and puts it on the signal', async () => { + const out = await admin.publishAdminState(STREAM, { patch: { pins: [pinOf(100)] } }); + + expect(published()).toHaveLength(1); + expect(published()[0]).toMatchObject({ type: 'ADMIN_STATE', rev: 8 }); + expect(out.messages).toHaveLength(1); + expect((await signalled()).snapshot).toMatchObject({ type: 'ADMIN_STATE', rev: 8 }); + }); + + it('splits a snapshot past the budget into a run every row of which fits', async () => { + const out = await admin.publishAdminState(STREAM, { patch: { pins: [pinOf(300 * 1024)] } }); + + const rows = published(); + const budget = CONFIG.media.imagePayloadMaxBytes - CONFIG.media.imagePayloadSafetyMarginBytes; + expect(rows.length).toBeGreaterThan(2); + expect(rows.at(-1)).toMatchObject({ type: 'admin_manifest', rev: 8, chunkCount: rows.length - 1 }); + for (const row of rows) expect(inflate(row)).toBeLessThanOrEqual(budget); + const [joined] = joinFramed(rows, ADMIN_FRAME); + expect(joined.payload).toMatchObject({ type: 'ADMIN_STATE', rev: 8 }); + expect(joined.payload.state.pins[0].snapshot.text).toHaveLength(300 * 1024); + expect(out.messages).toHaveLength(rows.length); + expect(out.published).toBe(out.messages.at(-1)); + }); + + it('signals a split snapshot by its rev only', async () => { + await admin.publishAdminState(STREAM, { patch: { pins: [pinOf(300 * 1024)] } }); + + const signal = await signalled(); + expect(signal).toEqual({ type: 'admin_invalidate', rev: 8, ts: expect.any(Number) }); + }); + + it('leaves the snapshot off a signal that would not fit, though the -3 row did', async () => { + streamrController.channelWireBytes.mockResolvedValueOnce(Number.MAX_SAFE_INTEGER); + + await admin.publishAdminState(STREAM, { patch: { pins: [pinOf(100)] } }); + + expect(published()).toHaveLength(1); + expect((await signalled()).snapshot).toBeUndefined(); + }); + + it('refuses a snapshot past the cap: nothing published, nothing applied', async () => { + const chars = (CONFIG.subscriptions.adminStateMaxChunks + 1) * 160 * 1024; + + await expect(admin.publishAdminState(STREAM, { patch: { pins: [pinOf(chars)] } })) + .rejects.toMatchObject({ code: 'ADMIN_STATE_TOO_LARGE' }); + + expect(streamrController.publishAdminState).not.toHaveBeenCalled(); + expect(streamrController.publishControl).not.toHaveBeenCalled(); + expect(channel.adminRev).toBe(7); + expect(manager.adminConfirm.track).not.toHaveBeenCalled(); + }); + + it('tracks the rev and ts the run carries, like a whole snapshot', async () => { + await admin.publishAdminState(STREAM, { patch: { pins: [pinOf(300 * 1024)] } }); + + const manifest = published().at(-1); + expect(manager.adminConfirm.track).toHaveBeenCalledWith(STREAM, + expect.objectContaining({ rev: 8, ts: manifest.ts })); + }); +}); + +describe('reading the -3 after a snapshot-less signal', () => { + afterEach(() => { + vi.useRealTimers(); + adminStatePoller.stop(); + }); + + it('polls once, after storage has had time, however many signals arrive', async () => { + vi.useFakeTimers(); + const admin = new AdminState({ channels: new Map() }); + const refresh = vi.fn(async () => {}); + adminStatePoller.start(STREAM, refresh); + refresh.mockClear(); + + admin.readAfterSignal(STREAM); + admin.readAfterSignal(STREAM); + await vi.advanceTimersByTimeAsync(CONFIG.subscriptions.adminSignalReadDelayMs - 1); + expect(refresh).not.toHaveBeenCalled(); + await vi.advanceTimersByTimeAsync(100); + + expect(refresh).toHaveBeenCalledTimes(1); + }); + + it('leaves a channel that is no longer being polled alone', async () => { + vi.useFakeTimers(); + const admin = new AdminState({ channels: new Map() }); + const refresh = vi.fn(async () => {}); + adminStatePoller.start(`${OWNER}/other-1`, refresh); + refresh.mockClear(); + + admin.readAfterSignal(STREAM); + await vi.advanceTimersByTimeAsync(CONFIG.subscriptions.adminSignalReadDelayMs + 100); + + expect(refresh).not.toHaveBeenCalled(); + }); +}); diff --git a/tests/unit/adminStateConfirm.test.js b/tests/unit/adminStateConfirm.test.js index 3c9970e..8dfd3da 100644 --- a/tests/unit/adminStateConfirm.test.js +++ b/tests/unit/adminStateConfirm.test.js @@ -106,6 +106,17 @@ describe('AdminStateConfirm', () => { expect(resendAdminState).toHaveBeenCalledTimes((limit + 1) * CONFIG.subscriptions.adminConfirmDelaysMs.length); }); + it('says the snapshot is too large rather than promising a retry that cannot land', async () => { + resendAdminState.mockResolvedValue(null); + manager.publishAdminState.mockRejectedValue( + Object.assign(new Error('too large'), { code: 'ADMIN_STATE_TOO_LARGE' })); + confirm.track(STREAM, { rev: 1, ts: 100 }); + await settle(); + expect(confirm.pending(STREAM)).toBeNull(); + expect(manager.notifyHandlers).toHaveBeenCalledWith('admin_state_too_large', { streamId: STREAM, rev: 1 }); + expect(manager.notifyHandlers).not.toHaveBeenCalledWith('admin_state_unconfirmed', expect.anything()); + }); + it('adopts a newer snapshot published from another device instead of republishing over it', async () => { const theirs = onStorage(5, 999); resendAdminState.mockResolvedValue(theirs); diff --git a/tests/unit/messageFlow.adminInvalidate.test.js b/tests/unit/messageFlow.adminInvalidate.test.js new file mode 100644 index 0000000..307bb45 --- /dev/null +++ b/tests/unit/messageFlow.adminInvalidate.test.js @@ -0,0 +1,62 @@ +/** + * The admin_invalidate signal on the -2: with the snapshot aboard it is + * applied inline; without it (the snapshot went out split on the -3) the + * receiver reads the -3, and never takes a snapshot from anyone but the owner. + */ + +import { describe, it, expect, beforeEach, vi } from 'vitest'; + +vi.mock('../../src/js/logger.js', () => ({ Logger: { debug: vi.fn(), info: vi.fn(), warn: vi.fn(), error: vi.fn() } })); +vi.mock('../../src/js/streamr.js', () => ({ streamrController: {}, STREAM_CONFIG: { MESSAGE_STREAM: { MESSAGES: 0, CONTROL: 1, MODERATION: 2 } } })); +vi.mock('../../src/js/auth.js', () => ({ authManager: { isConnected: () => true, getAddress: () => '0xmember' } })); +vi.mock('../../src/js/identity.js', () => ({ identityManager: {} })); +vi.mock('../../src/js/secureStorage.js', () => ({ secureStorage: {} })); +vi.mock('../../src/js/relayManager.js', () => ({ relayManager: {} })); +vi.mock('../../src/js/dm.js', () => ({ dmManager: {} })); +vi.mock('../../src/js/media.js', () => ({ mediaController: {} })); +vi.mock('../../src/js/adminStatePoller.js', () => ({ adminStatePoller: { getStreamId: () => null, markFresh: vi.fn() } })); + +const { MessageFlow } = await import('../../src/js/channels/MessageFlow.js'); + +const OWNER = '0xowner'; +const STREAM = `${OWNER}/chan-1`; +const snapshot = { type: 'ADMIN_STATE', rev: 8, ts: 800, createdBy: OWNER, state: {} }; + +let flow, manager, channel; + +beforeEach(() => { + channel = { createdBy: OWNER, adminRev: 7 }; + manager = { + channels: new Map([[STREAM, channel]]), + _isValidAdminState: (m) => !!m && m.type === 'ADMIN_STATE' && typeof m.rev === 'number' && !!m.state, + handleAdminMessage: vi.fn(), + adminState: { readAfterSignal: vi.fn() } + }; + flow = new MessageFlow(manager); +}); + +describe('admin_invalidate', () => { + it('applies a snapshot that rode along', async () => { + await flow.handleControlMessage(STREAM, { type: 'admin_invalidate', rev: 8, ts: 800, snapshot, account: OWNER }); + expect(manager.handleAdminMessage).toHaveBeenCalledWith(STREAM, snapshot); + expect(manager.adminState.readAfterSignal).not.toHaveBeenCalled(); + }); + + it('reads the -3 for a newer snapshot that did not ride along', async () => { + await flow.handleControlMessage(STREAM, { type: 'admin_invalidate', rev: 8, ts: 800, account: OWNER }); + expect(manager.adminState.readAfterSignal).toHaveBeenCalledWith(STREAM); + expect(manager.handleAdminMessage).not.toHaveBeenCalled(); + }); + + it('ignores a signal for a rev it already has', async () => { + await flow.handleControlMessage(STREAM, { type: 'admin_invalidate', rev: 7, ts: 700, account: OWNER }); + expect(manager.adminState.readAfterSignal).not.toHaveBeenCalled(); + }); + + it('ignores a signal from someone other than the owner', async () => { + await flow.handleControlMessage(STREAM, { type: 'admin_invalidate', rev: 8, ts: 800, account: '0xstranger' }); + expect(manager.adminState.readAfterSignal).not.toHaveBeenCalled(); + await flow.handleControlMessage(STREAM, { type: 'admin_invalidate', rev: 8, ts: 800, snapshot, account: '0xstranger' }); + expect(manager.handleAdminMessage).not.toHaveBeenCalled(); + }); +}); diff --git a/tests/unit/storageCopy.test.js b/tests/unit/storageCopy.test.js index 7337418..8f0ad77 100644 --- a/tests/unit/storageCopy.test.js +++ b/tests/unit/storageCopy.test.js @@ -66,6 +66,9 @@ let copy; let clock; let held; // rows the new provider holds: `${streamId}|${partition}|${ts}` +/** What publishAdminState returns for a snapshot that went out whole. */ +const published = (message) => ({ published: message, messages: [message] }); + beforeEach(() => { localStorage.clear(); [resendChannelImage, publishPasswordChallenge, probeStream, storedOn, republishAnchors].forEach((f) => f.mockReset()); @@ -87,7 +90,7 @@ beforeEach(() => { isChannelOwner: vi.fn(() => true), notifyHandlers: vi.fn(), bootstrapAdminState: vi.fn(async () => {}), - publishAdminState: vi.fn(async () => ({ published: { timestamp: ++clock, sequenceNumber: 0 } })), + publishAdminState: vi.fn(async () => published({ timestamp: ++clock, sequenceNumber: 0 })), ttlRepublish: { republishImage: vi.fn(async () => ({ timestamp: ++clock, sequenceNumber: 0 })) } }; republishAnchors.mockImplementation(async () => [ @@ -112,7 +115,7 @@ function providerStoresEverything() { manager.publishAdminState.mockImplementation(async () => { const ts = ++clock; held.add(`${ADMIN}|0|${ts}`); - return { published: { timestamp: ts, sequenceNumber: 0 } }; + return published({ timestamp: ts, sequenceNumber: 0 }); }); manager.ttlRepublish.republishImage.mockImplementation(async () => { const ts = ++clock; @@ -171,7 +174,7 @@ describe('StorageCopy.copyTo()', () => { const ts = ++clock; if (!adminLost) held.add(`${ADMIN}|0|${ts}`); adminLost = false; - return { published: { timestamp: ts, sequenceNumber: 0 } }; + return published({ timestamp: ts, sequenceNumber: 0 }); }); const outcome = await copy.copyTo(STREAM, NEW_NODE, null); @@ -181,6 +184,23 @@ describe('StorageCopy.copyTo()', () => { expect(republishAnchors).toHaveBeenCalledTimes(1); }); + it('confirms every row of a snapshot that went out split', async () => { + providerStoresEverything(); + let firstRun = true; + manager.publishAdminState.mockImplementation(async () => { + const rows = [0, 1, 2].map(() => ({ timestamp: ++clock, sequenceNumber: 0 })); + // The first run loses its middle chunk; only the manifest landing is not enough. + rows.forEach((r, i) => { if (!firstRun || i !== 1) held.add(`${ADMIN}|0|${r.timestamp}`); }); + firstRun = false; + return { published: rows.at(-1), messages: rows }; + }); + + const outcome = await copy.copyTo(STREAM, NEW_NODE, null); + + expect(outcome).toBe('present'); + expect(manager.publishAdminState).toHaveBeenCalledTimes(2); + }); + it('keeps the copy pending, and says so, when the provider never holds it', async () => { const outcome = await copy.copyTo(STREAM, NEW_NODE, null); diff --git a/tests/unit/streamr.adminRuns.test.js b/tests/unit/streamr.adminRuns.test.js new file mode 100644 index 0000000..9a12269 --- /dev/null +++ b/tests/unit/streamr.adminRuns.test.js @@ -0,0 +1,148 @@ +/** + * Reading an ADMIN_STATE that went out split. + * + * A snapshot too big for one wire message travels on the -3 as a run of + * chunks closed by a manifest. The reader has to put it back together from + * whatever window it read, never from rows of someone else, and never apply + * a run that is missing a row. + */ + +import { describe, it, expect, beforeEach, vi } from 'vitest'; + +vi.mock('../../src/js/auth.js', () => ({ + authManager: { getSigner: vi.fn() } +})); + +// Plain fixtures carry no signature; authority is resolveAuthor's, stubbed below. +vi.mock('../../src/js/envelopeSigner.js', async (importOriginal) => ({ + ...(await importOriginal()), + verifyEnvelopeAuthenticity: () => true, +})); + +import { streamrController, STREAM_CONFIG } from '../../src/js/streamr.js'; +import { splitFramed, ADMIN_FRAME } from '../../src/js/syncChunks.js'; + +const OWNER = '0xowner'; +const ADMIN = '0xowner/chan-3'; + +function iterator(messages) { + let i = 0; + return { + [Symbol.asyncIterator]() { + return { next: async () => (i < messages.length ? { done: false, value: messages[i++] } : { done: true }) }; + } + }; +} + +const row = (content, publisherId = OWNER) => ({ + content, publisherId, timestamp: 1, getPublisherId: () => publisherId +}); + +const snapshot = (rev, ts, pinText = '') => ({ + type: 'ADMIN_STATE', v: 1, rev, ts, createdBy: OWNER, + state: { bannedMembers: [], hiddenMessageIds: ['m-1'], pins: [{ targetId: 'm-9', snapshot: { text: pinText } }], absorbedThrough: 0 } +}); + +/** A snapshot split into `chunks` chunks plus its manifest, oldest row first. */ +const run = (s, chunks, runId = 'r1') => { + const size = JSON.stringify(s).length; + return splitFramed(s, runId, ADMIN_FRAME, Math.ceil(size / chunks)); +}; + +describe('resendAdminState with split snapshots', () => { + let windows; // what each resend returns, newest rows last, by `last` + + beforeEach(() => { + windows = null; + streamrController.client = { + resend: vi.fn(async (_part, { last }) => iterator(windows(last))) + }; + vi.spyOn(streamrController, '_gatedChannelFor').mockResolvedValue(null); + vi.spyOn(streamrController, 'publisherMayWrite').mockResolvedValue(true); + vi.spyOn(streamrController, 'resolveAuthor').mockImplementation(async (_s, _m, p) => p); + }); + + /** The -3 as a list of rows, read `last` at a time from the end. */ + const storage = (rows) => (last) => rows.slice(-last); + + it('joins a run into the snapshot it carried', async () => { + const big = snapshot(5, 500, 'x'.repeat(900)); + windows = storage([row(snapshot(4, 400)), ...run(big, 3).map(r => row(r))]); + + const latest = await streamrController.resendAdminState(ADMIN, { historyCount: 10 }); + + expect(latest).toEqual(big); + }); + + it('ranks a run against whole snapshots by rev', async () => { + const big = snapshot(5, 500, 'x'.repeat(900)); + windows = storage([...run(big, 2).map(r => row(r)), row(snapshot(6, 600))]); + + const latest = await streamrController.resendAdminState(ADMIN, { historyCount: 10 }); + + expect(latest.rev).toBe(6); + }); + + it('never completes a run with rows from another publisher', async () => { + const big = snapshot(5, 500, 'x'.repeat(900)); + const rows = run(big, 3); + windows = storage([ + row(snapshot(4, 400)), + ...rows.map((r, i) => row(r, i === 1 ? '0xintruder' : OWNER)) + ]); + + const latest = await streamrController.resendAdminState(ADMIN, { historyCount: 10 }); + + expect(latest.rev).toBe(4); + }); + + it('reads again, wider, when the newest run reaches past the window', async () => { + const big = snapshot(5, 500, 'x'.repeat(900)); + windows = storage([row(snapshot(4, 400)), ...run(big, 6).map(r => row(r))]); + + const latest = await streamrController.resendAdminState(ADMIN, { historyCount: 5 }); + + expect(latest).toEqual(big); + expect(streamrController.client.resend).toHaveBeenCalledTimes(2); + expect(streamrController.client.resend.mock.calls[1][1]).toEqual({ last: 5 + 6 + 1, raw: true }); + }); + + it('does not read again for a run the window could hold: a lost row stays lost', async () => { + const big = snapshot(5, 500, 'x'.repeat(900)); + const rows = run(big, 3); + windows = storage([row(snapshot(4, 400)), ...rows.filter((_, i) => i !== 1).map(r => row(r))]); + + const latest = await streamrController.resendAdminState(ADMIN, { historyCount: 10 }); + + expect(latest.rev).toBe(4); + expect(streamrController.client.resend).toHaveBeenCalledTimes(1); + }); + + it('does not chase a cut run older than the snapshot it already has', async () => { + const big = snapshot(5, 500, 'x'.repeat(900)); + windows = storage([...run(big, 6).map(r => row(r)), row(snapshot(6, 600))]); + + const latest = await streamrController.resendAdminState(ADMIN, { historyCount: 5 }); + + expect(latest.rev).toBe(6); + expect(streamrController.client.resend).toHaveBeenCalledTimes(1); + }); + + it('drops a run whose rows disagree with their manifest', async () => { + const big = snapshot(5, 500, 'x'.repeat(900)); + const rows = run(big, 2); + rows.forEach(r => { r.rev = 9; }); + windows = storage([row(snapshot(4, 400)), ...rows.map(r => row(r))]); + + const latest = await streamrController.resendAdminState(ADMIN, { historyCount: 10 }); + + expect(latest.rev).toBe(4); + }); + + it('reads the moderation partition', async () => { + windows = storage([]); + await streamrController.resendAdminState(ADMIN, { historyCount: 5 }); + expect(streamrController.client.resend.mock.calls[0][0]) + .toEqual({ streamId: ADMIN, partition: STREAM_CONFIG.ADMIN_STREAM.MODERATION }); + }); +}); diff --git a/tests/unit/streamr.wireBytes.test.js b/tests/unit/streamr.wireBytes.test.js new file mode 100644 index 0000000..88002b9 --- /dev/null +++ b/tests/unit/streamr.wireBytes.test.js @@ -0,0 +1,134 @@ +/** + * What an ADMIN_STATE and a control signal weigh on the wire. + * + * Past the DataChannel's max-message-size the transport throws inside the + * SDK, where the publisher never sees it, so the snapshot is split by these + * numbers. They mirror publishAdminState and publishAsChannel branch for + * branch; what pins them is that measure and publish agree. + */ + +import { describe, it, expect, beforeEach, beforeAll, vi } from 'vitest'; +import { ethers } from 'ethers'; + +globalThis.ethers = ethers; + +const ACCOUNT_KEY = '0x' + '11'.repeat(32); + +vi.mock('../../src/js/auth.js', () => ({ + authManager: { + wallet: { privateKey: '0x' + '11'.repeat(32) }, + getAddress: () => new ethers.Wallet('0x' + '11'.repeat(32)).address, + isConnected: () => true + } +})); + +const { streamrController } = await import('../../src/js/streamr.js'); +const { cryptoManager } = await import('../../src/js/crypto.js'); +const { epochKeyCrypto } = await import('../../src/js/epochKeyCrypto.js'); +const { epochKeyManager } = await import('../../src/js/epochKeyManager.js'); +const { channelManager } = await import('../../src/js/channels.js'); +const { clearChannelIdentities } = await import('../../src/js/channelIdentity.js'); + +const snapshot = { + type: 'ADMIN_STATE', v: 1, rev: 3, ts: 1789000000000, createdBy: '0xowner', + state: { bannedMembers: [], hiddenMessageIds: ['m-1'], pins: [{ targetId: 'm-2', snapshot: { text: 'olá "pin" 🐦' } }] } +}; + +const bytes = (content) => new TextEncoder().encode(JSON.stringify(content)).length; + +describe('encryptedLength', () => { + it('is the length encrypt() returns, without deriving a key', async () => { + for (const text of ['', 'a', 'ab', 'abc', 'olá 🐦 "x"', 'y'.repeat(4097)]) { + const out = await cryptoManager.encrypt(text, 'pw'); + expect(cryptoManager.encryptedLength(new TextEncoder().encode(text).length)).toBe(out.length); + } + }); +}); + +describe('wire bytes of the -3 and the -2', () => { + let published; + + beforeAll(() => { + window.EthereumKeyPairIdentity = { + fromPrivateKey: (pk) => ({ + getUserId: async () => new ethers.Wallet(pk).address, + getSignatureType: () => 'ECDSA_SECP256K1_EVM' + }) + }; + window.SignatureType = { ERC_1271: 3, ECDSA_SECP256K1_EVM: 2 }; + }); + + beforeEach(() => { + published = []; + clearChannelIdentities(); + channelManager.channels.clear(); + streamrController.client = { + publish: vi.fn(async (_part, content) => { published.push(content); return { timestamp: 1 }; }) + }; + vi.spyOn(streamrController, 'publishAs').mockImplementation( + async (_identity, _streamId, _partition, content) => { published.push(content); return { timestamp: 1 }; }); + streamrController._accountIdentity = window.EthereumKeyPairIdentity.fromPrivateKey(ACCOUNT_KEY); + }); + + describe('public channel', () => { + it('measures the -3 publish', async () => { + const measured = await streamrController.adminWireBytes('0xowner/pub-3', snapshot); + await streamrController.publishAdminState('0xowner/pub-3', snapshot); + expect(measured).toBe(bytes(published[0])); + }); + + it('measures the -2 publish, proof included', async () => { + const signal = { type: 'admin_invalidate', rev: 3, ts: 1, snapshot }; + const measured = await streamrController.channelWireBytes('0xowner/pub-2', signal); + await streamrController.publishAsChannel('0xowner/pub-2', 0, signal); + expect(measured).toBe(bytes(published[0])); + }); + }); + + describe('password channel', () => { + it('measures the -3 publish', async () => { + const measured = await streamrController.adminWireBytes('0xowner/pw-3', snapshot, 'pw'); + await streamrController.publishAdminState('0xowner/pw-3', snapshot, 'pw'); + expect(measured).toBe(bytes(published[0])); + }); + + it('measures the -2 publish', async () => { + const signal = { type: 'admin_invalidate', rev: 3, ts: 1, snapshot }; + const measured = await streamrController.channelWireBytes('0xowner/pw-2', signal, 'pw'); + await streamrController.publishAsChannel('0xowner/pw-2', 0, signal, 'pw'); + expect(measured).toBe(bytes(published[0])); + }); + }); + + describe('gated channel', () => { + const channel = { messageStreamId: '0xowner/gated-1', type: 'gated', gate: { address: '0xgate' } }; + + beforeEach(async () => { + const key = { kid: '7.abcdef012345', cryptoKey: await epochKeyCrypto.importEpochKey(epochKeyCrypto.generateEpochKey()) }; + vi.spyOn(epochKeyManager, 'getCurrentKey').mockResolvedValue(key); + channelManager.channels.set(channel.messageStreamId, channel); + }); + + it('measures the -3 publish', async () => { + const measured = await streamrController.adminWireBytes('0xowner/gated-3', snapshot); + await streamrController.publishAdminState('0xowner/gated-3', snapshot); + expect(measured).toBe(bytes(published[0])); + }); + + it('measures the -2 publish', async () => { + const signal = { type: 'admin_invalidate', rev: 3, ts: 1, snapshot }; + const measured = await streamrController.channelWireBytes('0xowner/gated-2', signal); + await streamrController.publishAsChannel('0xowner/gated-2', 0, signal); + expect(measured).toBe(bytes(published[0])); + }); + + it('tries to recover a missing key before sizing, as the publish does', async () => { + epochKeyManager.getCurrentKey.mockResolvedValueOnce(null); + const ensure = vi.spyOn(epochKeyManager, 'ensureChannelKeys').mockResolvedValue(); + + await streamrController.adminWireBytes('0xowner/gated-3', snapshot); + + expect(ensure).toHaveBeenCalledWith(channel); + }); + }); +}); From 6d4ce7b3509630c46f2cf659fdb6c9d5cdd78e13 Mon Sep 17 00:00:00 2001 From: Ocnrb Date: Fri, 25 Sep 2026 08:44:02 +0100 Subject: [PATCH 3/4] Read the admin stream again when an announced revision has not landed yet Co-authored-by: Claude Opus 5.5 --- src/js/channels/AdminState.js | 30 ++++++--- src/js/channels/MessageFlow.js | 2 +- src/js/config.js | 7 ++- src/js/subscriptionManager.js | 2 +- tests/unit/adminState.split.test.js | 62 +++++++++++++++---- .../unit/messageFlow.adminInvalidate.test.js | 2 +- 6 files changed, 79 insertions(+), 26 deletions(-) diff --git a/src/js/channels/AdminState.js b/src/js/channels/AdminState.js index 9257962..1b3c95d 100644 --- a/src/js/channels/AdminState.js +++ b/src/js/channels/AdminState.js @@ -480,15 +480,31 @@ export class AdminState { /** * An admin_invalidate announced a snapshot too big to ride along: poll - * the -3 once storage has had time to hold it. Signals arriving while a - * read is pending share it. + * 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) { + readAfterSignal(messageStreamId, rev) { if (this._signalReads.has(messageStreamId)) return; - this._signalReads.set(messageStreamId, setTimeout(() => { - this._signalReads.delete(messageStreamId); - if (adminStatePoller.getStreamId() === messageStreamId) adminStatePoller.pollNow(); - }, CONFIG.subscriptions.adminSignalReadDelayMs)); + 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 ---------- diff --git a/src/js/channels/MessageFlow.js b/src/js/channels/MessageFlow.js index a58a305..dc11e60 100644 --- a/src/js/channels/MessageFlow.js +++ b/src/js/channels/MessageFlow.js @@ -83,7 +83,7 @@ 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); + this.manager.adminState?.readAfterSignal(streamId, incomingRev); return; } if (!this.manager._isValidAdminState(data.snapshot)) return; diff --git a/src/js/config.js b/src/js/config.js index cb067f9..ed7f977 100644 --- a/src/js/config.js +++ b/src/js/config.js @@ -298,9 +298,10 @@ export const CONFIG = { // 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 this late, - // once storage has had time to hold the run. - adminSignalReadDelayMs: 5000, + // 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], diff --git a/src/js/subscriptionManager.js b/src/js/subscriptionManager.js index 94c36c0..0533a10 100644 --- a/src/js/subscriptionManager.js +++ b/src/js/subscriptionManager.js @@ -587,7 +587,7 @@ class SubscriptionManager { // published by the same admin. We additionally inject createdBy // from the account so applyAdminState's owner check passes. if (msg.snapshot === undefined) { - channelManager.adminState?.readAfterSignal(streamId); + channelManager.adminState?.readAfterSignal(streamId, typeof msg.rev === 'number' ? msg.rev : 0); return; } if (!msg.snapshot || typeof msg.snapshot !== 'object') return; diff --git a/tests/unit/adminState.split.test.js b/tests/unit/adminState.split.test.js index 0e13ac4..8d12927 100644 --- a/tests/unit/adminState.split.test.js +++ b/tests/unit/adminState.split.test.js @@ -136,36 +136,72 @@ describe('ADMIN_STATE publish on the wire', () => { }); describe('reading the -3 after a snapshot-less signal', () => { + const pollInterval = CONFIG.subscriptions.adminPollIntervalMs; + + // The regular poll is not what is under test: keep its ticks out of the count. + beforeEach(() => { CONFIG.subscriptions.adminPollIntervalMs = 24 * 3600 * 1000; }); + afterEach(() => { vi.useRealTimers(); adminStatePoller.stop(); + CONFIG.subscriptions.adminPollIntervalMs = pollInterval; }); - it('polls once, after storage has had time, however many signals arrive', async () => { - vi.useFakeTimers(); - const admin = new AdminState({ channels: new Map() }); + const [FIRST, SECOND] = CONFIG.subscriptions.adminSignalReadDelaysMs; + + /** An AdminState over one channel at rev 7, polled with `refresh`. */ + const polled = (streamId = STREAM) => { + const channel = { messageStreamId: STREAM, adminRev: 7 }; + const admin = new AdminState({ channels: new Map([[STREAM, channel]]) }); const refresh = vi.fn(async () => {}); - adminStatePoller.start(STREAM, refresh); + adminStatePoller.start(streamId, refresh); refresh.mockClear(); + return { admin, channel, refresh }; + }; - admin.readAfterSignal(STREAM); - admin.readAfterSignal(STREAM); - await vi.advanceTimersByTimeAsync(CONFIG.subscriptions.adminSignalReadDelayMs - 1); + it('polls after storage has had time, however many signals arrive', async () => { + vi.useFakeTimers(); + const { admin, refresh } = polled(); + + admin.readAfterSignal(STREAM, 8); + admin.readAfterSignal(STREAM, 8); + await vi.advanceTimersByTimeAsync(FIRST - 1); expect(refresh).not.toHaveBeenCalled(); await vi.advanceTimersByTimeAsync(100); expect(refresh).toHaveBeenCalledTimes(1); }); + it('reads once more while the announced rev has not landed, and never again', async () => { + vi.useFakeTimers(); + const { admin, refresh } = polled(); + + admin.readAfterSignal(STREAM, 8); + await vi.advanceTimersByTimeAsync(FIRST + SECOND + 100); + expect(refresh).toHaveBeenCalledTimes(2); + await vi.advanceTimersByTimeAsync(10 * (FIRST + SECOND)); + + expect(refresh).toHaveBeenCalledTimes(2); + }); + + it('does not read again once the announced rev has landed', async () => { + vi.useFakeTimers(); + const { admin, channel, refresh } = polled(); + + admin.readAfterSignal(STREAM, 8); + await vi.advanceTimersByTimeAsync(FIRST + 100); + channel.adminRev = 8; + await vi.advanceTimersByTimeAsync(SECOND + 100); + + expect(refresh).toHaveBeenCalledTimes(1); + }); + it('leaves a channel that is no longer being polled alone', async () => { vi.useFakeTimers(); - const admin = new AdminState({ channels: new Map() }); - const refresh = vi.fn(async () => {}); - adminStatePoller.start(`${OWNER}/other-1`, refresh); - refresh.mockClear(); + const { admin, refresh } = polled(`${OWNER}/other-1`); - admin.readAfterSignal(STREAM); - await vi.advanceTimersByTimeAsync(CONFIG.subscriptions.adminSignalReadDelayMs + 100); + admin.readAfterSignal(STREAM, 8); + await vi.advanceTimersByTimeAsync(FIRST + SECOND + 100); expect(refresh).not.toHaveBeenCalled(); }); diff --git a/tests/unit/messageFlow.adminInvalidate.test.js b/tests/unit/messageFlow.adminInvalidate.test.js index 307bb45..9820bc3 100644 --- a/tests/unit/messageFlow.adminInvalidate.test.js +++ b/tests/unit/messageFlow.adminInvalidate.test.js @@ -44,7 +44,7 @@ describe('admin_invalidate', () => { it('reads the -3 for a newer snapshot that did not ride along', async () => { await flow.handleControlMessage(STREAM, { type: 'admin_invalidate', rev: 8, ts: 800, account: OWNER }); - expect(manager.adminState.readAfterSignal).toHaveBeenCalledWith(STREAM); + expect(manager.adminState.readAfterSignal).toHaveBeenCalledWith(STREAM, 8); expect(manager.handleAdminMessage).not.toHaveBeenCalled(); }); From dcc77616823db42e075b54ef56bb0d2c4bb8b24a Mon Sep 17 00:00:00 2001 From: Ocnrb Date: Fri, 25 Sep 2026 08:51:45 +0100 Subject: [PATCH 4/4] Probe the newest admin row before reading the window on each poll Co-authored-by: Claude Opus 5.5 --- src/js/channels/AdminState.js | 11 ++++ src/js/streamr.js | 32 ++++++++-- tests/unit/adminState.probe.test.js | 92 ++++++++++++++++++++++++++++ tests/unit/streamr.adminRuns.test.js | 38 ++++++++++++ 4 files changed, 169 insertions(+), 4 deletions(-) create mode 100644 tests/unit/adminState.probe.test.js diff --git a/src/js/channels/AdminState.js b/src/js/channels/AdminState.js index 1b3c95d..c523844 100644 --- a/src/js/channels/AdminState.js +++ b/src/js/channels/AdminState.js @@ -281,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. diff --git a/src/js/streamr.js b/src/js/streamr.js index 7f76e6d..c0a1496 100644 --- a/src/js/streamr.js +++ b/src/js/streamr.js @@ -3560,17 +3560,40 @@ class StreamrController { return read.latest; } + /** + * The newest -3/P0 row alone, through the same opening and authority + * checks as the window: `{ type, rev, ts, publisherId }` when it is a + * whole snapshot or a run's manifest, null for anything else and on any + * error. + * @param {string} adminStreamId - Admin stream ID (ends with -3) + * @param {Object} [options] + * @param {string|null} [options.password=null] Channel password (encrypted channels) + * @returns {Promise} + */ + async probeAdminState(adminStreamId, { password = null } = {}) { + if (!this.client) return null; + try { + const { newest } = await this._readAdminWindow(adminStreamId, 1, password); + return newest && (newest.type === 'ADMIN_STATE' || newest.type === ADMIN_FRAME.manifest) + ? newest : null; + } catch (error) { + Logger.debug('probeAdminState error (window read follows):', error.message); + return null; + } + } + /** * One resend of the -3/P0 window: the newest snapshot among whole * ADMIN_STATE rows and complete runs, plus the manifest of a newer run - * the window cut short, if any. Rows only join a run of their own - * publisher, and each row passes the same authority check a whole - * snapshot does. + * the window cut short, if any, and the newest row that passed. Rows + * only join a run of their own publisher, and each row passes the same + * authority check a whole snapshot does. * @private */ async _readAdminWindow(adminStreamId, last, password) { const partition = STREAM_CONFIG.ADMIN_STREAM.MODERATION; let latest = null; + let newest = null; const framed = new Map(); // publisherId -> opened chunk/manifest rows const num = (v) => (typeof v === 'number' ? v : 0); const newer = (a, b) => !b @@ -3638,6 +3661,7 @@ class StreamrController { const publisherId = await this.resolveAuthor( adminStreamId, message, transportPublisher); if (!publisherId) continue; + newest = { type: content.type || 'ADMIN_STATE', rev: num(content.rev), ts: num(content.ts), publisherId }; if (isFramed) { if (!framed.has(publisherId)) framed.set(publisherId, []); @@ -3672,7 +3696,7 @@ class StreamrController { if (cutRun && (!newer(cutRun, latest) || cutRun.chunkCount + 1 <= last || cutRun.chunkCount > CONFIG.subscriptions.adminReadMaxChunks)) cutRun = null; - return { latest, cutRun }; + return { latest, cutRun, newest }; } /** diff --git a/tests/unit/adminState.probe.test.js b/tests/unit/adminState.probe.test.js new file mode 100644 index 0000000..8766439 --- /dev/null +++ b/tests/unit/adminState.probe.test.js @@ -0,0 +1,92 @@ +/** + * The poll of an open channel's moderation reads the newest -3 row first. + * Only an owner's snapshot or manifest at a rev already held ends the poll + * there; anything else reads the window exactly as before, so a probe can + * save a read but never hide a change. + */ + +import { describe, it, expect, beforeEach, vi } from 'vitest'; + +vi.mock('../../src/js/logger.js', () => ({ + Logger: { debug: vi.fn(), info: vi.fn(), warn: vi.fn(), error: vi.fn() } +})); + +const OWNER = '0x' + 'aa'.repeat(20); +const STREAM = `${OWNER}/room-1`; + +vi.mock('../../src/js/auth.js', () => ({ authManager: { getAddress: () => '0x' + 'bb'.repeat(20) } })); +vi.mock('../../src/js/streamr.js', () => ({ + streamrController: { probeAdminState: vi.fn(), resendAdminState: vi.fn() }, + STREAM_CONFIG: { ADMIN_HISTORY_COUNT: 10 }, + deriveEphemeralId: (id) => id.replace(/-1$/, '-2'), + deriveAdminId: (id) => id.replace(/-1$/, '-3') +})); + +const { AdminState } = await import('../../src/js/channels/AdminState.js'); +const { streamrController } = await import('../../src/js/streamr.js'); + +const snapshot = (rev) => ({ + type: 'ADMIN_STATE', rev, ts: rev * 100, createdBy: OWNER, + state: { bannedMembers: [], hiddenMessageIds: [`m-${rev}`], pins: [] } +}); + +describe('refreshAdminState probe', () => { + let admin; + let channel; + let manager; + + beforeEach(() => { + vi.clearAllMocks(); + channel = { messageStreamId: STREAM, createdBy: OWNER, adminLoaded: true, adminRev: 5, adminTs: 500 }; + manager = { channels: new Map([[STREAM, channel]]), notifyHandlers: vi.fn() }; + admin = new AdminState(manager); + for (const m of ['_isValidAdminState', '_normalizeAdminState', 'applyAdminState']) manager[m] = admin[m].bind(admin); + streamrController.resendAdminState.mockResolvedValue(snapshot(6)); + }); + + const probe = (row) => streamrController.probeAdminState.mockResolvedValue(row); + + it("stops at the owner's snapshot at a rev already held", async () => { + probe({ type: 'ADMIN_STATE', rev: 5, ts: 500, publisherId: OWNER }); + + expect(await admin.refreshAdminState(STREAM)).toBe(false); + expect(streamrController.resendAdminState).not.toHaveBeenCalled(); + }); + + it("stops at the owner's manifest at a rev already held", async () => { + probe({ type: 'admin_manifest', rev: 5, ts: 500, publisherId: OWNER }); + + expect(await admin.refreshAdminState(STREAM)).toBe(false); + expect(streamrController.resendAdminState).not.toHaveBeenCalled(); + }); + + it('reads the window for a newer rev', async () => { + probe({ type: 'admin_manifest', rev: 6, ts: 600, publisherId: OWNER }); + + expect(await admin.refreshAdminState(STREAM)).toBe(true); + expect(streamrController.resendAdminState).toHaveBeenCalledWith(`${OWNER}/room-3`, { historyCount: 5, password: null }); + expect(channel.adminRev).toBe(6); + }); + + it('reads the window when the newest row is not the owner\'s', async () => { + probe({ type: 'ADMIN_STATE', rev: 5, ts: 500, publisherId: '0x' + 'cc'.repeat(20) }); + + expect(await admin.refreshAdminState(STREAM)).toBe(true); + expect(streamrController.resendAdminState).toHaveBeenCalledTimes(1); + }); + + it('reads the window when the probe has nothing to say', async () => { + probe(null); + + expect(await admin.refreshAdminState(STREAM)).toBe(true); + expect(streamrController.resendAdminState).toHaveBeenCalledTimes(1); + }); + + it('takes the owner from the stream namespace when the record has no createdBy', async () => { + delete channel.createdBy; + probe({ type: 'ADMIN_STATE', rev: 5, ts: 500, publisherId: OWNER }); + + expect(await admin.refreshAdminState(STREAM)).toBe(false); + expect(streamrController.resendAdminState).not.toHaveBeenCalled(); + }); +}); diff --git a/tests/unit/streamr.adminRuns.test.js b/tests/unit/streamr.adminRuns.test.js index 9a12269..88c821f 100644 --- a/tests/unit/streamr.adminRuns.test.js +++ b/tests/unit/streamr.adminRuns.test.js @@ -139,6 +139,44 @@ describe('resendAdminState with split snapshots', () => { expect(latest.rev).toBe(4); }); + describe('probeAdminState', () => { + it('names the newest row when it is a whole snapshot', async () => { + windows = storage([row(snapshot(4, 400)), row(snapshot(5, 500))]); + + expect(await streamrController.probeAdminState(ADMIN)) + .toEqual({ type: 'ADMIN_STATE', rev: 5, ts: 500, publisherId: OWNER }); + expect(streamrController.client.resend.mock.calls[0][1]).toEqual({ last: 1, raw: true }); + }); + + it("names the newest row when it is a run's manifest", async () => { + windows = storage(run(snapshot(5, 500, 'x'.repeat(900)), 3).map(r => row(r))); + + expect(await streamrController.probeAdminState(ADMIN)) + .toEqual({ type: 'admin_manifest', rev: 5, ts: 500, publisherId: OWNER }); + }); + + it('has nothing to say about a chunk', async () => { + const rows = run(snapshot(5, 500, 'x'.repeat(900)), 3); + windows = storage(rows.slice(0, -1).map(r => row(r))); + + expect(await streamrController.probeAdminState(ADMIN)).toBeNull(); + }); + + it('has nothing to say about a row the author check drops', async () => { + // An old gated -3 where the clone still publishes: the signer is not the admin. + streamrController.resolveAuthor.mockResolvedValue(null); + windows = storage([row(snapshot(5, 500), '0xclone')]); + + expect(await streamrController.probeAdminState(ADMIN)).toBeNull(); + }); + + it('has nothing to say when the read fails', async () => { + streamrController.client.resend = vi.fn(async () => { throw new Error('network'); }); + + expect(await streamrController.probeAdminState(ADMIN)).toBeNull(); + }); + }); + it('reads the moderation partition', async () => { windows = storage([]); await streamrController.resendAdminState(ADMIN, { historyCount: 5 });