From b96cb8e15195b9b5991977a7022baa3541f1f62f Mon Sep 17 00:00:00 2001 From: Dagur Valberg Johannsson Date: Tue, 21 Apr 2026 14:50:54 +0200 Subject: [PATCH] Add chunk transport extension for oversized messages Splits ProtocolMessages exceeding NIP-44's 65,535-byte plaintext ceiling across multiple gift-wrapped events. Symmetric (both directions), fire-and-forget, backward-compatible via a new transport-level \`extensions\` field on \`dapp_ready\` and \`wallet_ready\`. Resolves the \"Failed to swap: invalid plaintext size\" error on aggregated swap requests and enables signed-tx responses up to the 1 MB BCH consensus limit (~2 MB hex). --- CLAUDE.md | 3 +- docs/extensions.md | 6 + docs/protocol.md | 3 + docs/transport.md | 132 +++++++ package-lock.json | 10 +- packages/core/package.json | 2 +- packages/core/src/index.ts | 6 + packages/core/src/protocols/base.ts | 8 + packages/core/src/protocols/hdwalletv1.ts | 35 ++ packages/core/src/relay-client.test.ts | 85 +++++ packages/core/src/relay-client.ts | 216 ++++++++---- packages/core/src/transforms/chunk.test.ts | 322 ++++++++++++++++++ packages/core/src/transforms/chunk.ts | 305 +++++++++++++++++ packages/dapp/package.json | 2 +- .../dapp/src/dapp-connection-manager.test.ts | 1 + packages/dapp/src/dapp-connection-manager.ts | 15 + packages/react/package.json | 2 +- packages/test-cli/package.json | 2 +- packages/wallet/package.json | 2 +- packages/wallet/src/integration/chunk.test.ts | 190 +++++++++++ packages/wallet/src/integration/helpers.ts | 19 ++ .../wallet/src/wallet-connection-manager.ts | 14 + 22 files changed, 1309 insertions(+), 71 deletions(-) create mode 100644 packages/core/src/transforms/chunk.test.ts create mode 100644 packages/core/src/transforms/chunk.ts create mode 100644 packages/wallet/src/integration/chunk.test.ts diff --git a/CLAUDE.md b/CLAUDE.md index 91033cc..48a66b9 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -60,7 +60,8 @@ Protocol and architecture documentation lives in `docs/`. Keep it up to date whe | File | Update when… | |------|-------------| | `docs/protocol.md` | Protocol messages, handshake logic, or `PathName`/`PathXpub`/`NextIndex` types change | -| `docs/extensions.md` | Extension system, `WalletAdapter` extension hooks, known extensions, or custom message conventions change | +| `docs/extensions.md` | hdwalletv1 protocol-level extensions, `WalletAdapter` extension hooks, known protocol extensions, or custom message conventions change | +| `docs/transport.md` (transport extensions section) | Transport-level extensions (`chunk`, future: compression), base-level `extensions` field on `dapp_ready`/`wallet_ready`, reassembly semantics | | `docs/connection-uri.md` | URI format, key exchange flow, or credential structure changes | | `docs/transport.md` | `RelayClient`, `initiateRelay`, reconnect logic, or encryption scheme changes | | `docs/wallet.md` | `WalletAdapter`, `WalletConnectionManager`, or connection lifecycle changes | diff --git a/docs/extensions.md b/docs/extensions.md index ac9720d..83f1166 100644 --- a/docs/extensions.md +++ b/docs/extensions.md @@ -1,5 +1,11 @@ # Extensions +This document covers `hdwalletv1` protocol-level extensions — optional capabilities that extend +the application protocol (extra path names, custom message actions, per-wallet features). For +transport-level capabilities that apply below the application protocol regardless of which +protocol is in use (chunking, future: compression), see +[transport.md § Transport-level extensions](transport.md#transport-level-extensions). + The hdwalletv1 protocol supports optional extensions that let wallets and dapps negotiate additional capabilities beyond the core sign-transaction flow. Extensions are backward-compatible: existing wallets and dapps that don't know about extensions continue to work unchanged. diff --git a/docs/protocol.md b/docs/protocol.md index 694fa6b..492c609 100644 --- a/docs/protocol.md +++ b/docs/protocol.md @@ -36,6 +36,9 @@ sign_transaction_request — dapp → wallet, asks wallet to sign a transaction sign_transaction_response — wallet → dapp, returns signed tx or error sign_cancel — dapp → wallet only, cancels an in-flight sign_transaction_request disconnect — either → either, courtesy notification before tearing down +chunk — either → either, transport-level. Carries one slice of a larger message + that exceeds NIP-44's 65,535-byte plaintext ceiling. Not tied to + hdwalletv1 semantics. See transport.md. ``` The action names above are the well-known set. Extensions may define additional action strings diff --git a/docs/transport.md b/docs/transport.md index 0e0ab63..d4e110e 100644 --- a/docs/transport.md +++ b/docs/transport.md @@ -223,3 +223,135 @@ nostr-tools `SimplePool.subscribeMany()` deduplicates events by ID — it tracks in a per-subscription `_knownIds` set and only fires `onevent` once per unique ID. Since `pool.publish(urls, event)` sends the identical event (same ID) to all relays, the receiving side's pool delivers it exactly once. + +## Transport-level extensions + +Transport-level extensions are capabilities of the relay/gift-wrap transport itself, independent +of any application protocol. They're distinct from the protocol-level extensions documented in +[extensions.md](extensions.md), which extend `hdwalletv1` specifically. + +### Negotiation + +Both sides advertise transport-level extensions in a base `extensions` field on their handshake +message (`dapp_ready` / `wallet_ready`). The shape is identical on both sides: + +```typescript +interface DappReadyMessage { + // ... + extensions?: Record; +} + +interface WalletReadyMessage { + // ... + extensions?: Record; +} +``` + +A capability is considered enabled only when **both** sides advertise it. `RelayClient` exposes +`setPeerCapabilities({...})` so the connection managers can inform it once they've parsed the +peer's `_ready` message. + +Extension values are per-extension; today `chunk` uses `{ version: 1 }`. Unknown extension keys +are ignored by implementations that don't recognise them, so adding new capabilities is always +backward-compatible — an old peer simply doesn't advertise them and the new side falls back to +the non-extended behaviour. + +### Known transport extensions + +| Extension | Purpose | +|---|---| +| `chunk` | Split messages larger than NIP-44's 65,535-byte plaintext ceiling across multiple gift-wrapped events. Symmetric — enabled when both sides advertise. See below. | + +## Chunking (`chunk` extension) + +### Why it exists + +NIP-44 caps plaintext at 65,535 bytes — the length is encoded as a U16BE prefix in the wire +format, so the limit is structural, not a guardrail that can be raised. NIP-17 gift-wrap +additionally encrypts twice (rumor inside seal inside wrap), so a single 50+ KB `ProtocolMessage` +will blow the outer wrap's plaintext budget even when the inner payload itself looks small enough. + +Real scenarios that exceed the cap: + +- Aggregated swap `sign_transaction_request` with many pool UTXOs (each with hex `lockingBytecode` + and `unlockingBytecode` in the per-input `sourceOutputs`). +- `sign_transaction_response` carrying a signed transaction hex string. Policy-max BCH transactions + are 100 KB (→ 200 KB hex); consensus-max is 1 MB (→ 2 MB hex). + +### Wire format + +```typescript +interface ChunkMessage extends ProtocolMessage { + action: "chunk"; + time: number; // shared across all chunks of one logical message + msgId: string; // random identifier, shared across all chunks + index: number; // 0-based chunk index within [0, total) + total: number; // total number of chunks for this msgId + data: string; // base64 slice of the UTF-8 bytes of JSON.stringify(originalMessage) +} +``` + +To reconstruct the original message: concatenate the `data` strings in `index` order, base64-decode +to bytes, UTF-8-decode to a string, `JSON.parse`. + +### Sender + +`RelayClient.publishMessage` measures the UTF-8 byte length of the serialized `ProtocolMessage`. +If it exceeds the per-message threshold: + +- If the peer's `chunk` capability is enabled, the message is split into `ChunkMessage`s and each + is published individually via the same `wrapEvent` path as any other message. No ACKs — see + "Failure modes" below. +- If the peer has not advertised `chunk`, `publishMessage` throws a structured error + (`"Cannot send : message is larger than NIP-44's 65,535-byte ceiling and the peer does + not advertise the 'chunk' transport extension. Please update the connected wallet/dapp..."`). + This replaces the cryptic `invalid plaintext size` error from nostr-tools. + +The per-chunk budget is sized conservatively. Each chunk's raw data is ≤ 30,000 bytes, which after +base64 expansion (~4/3×), JSON envelope overhead, and the two-layer NIP-17 gift-wrap encryption +stays well under NIP-44's 65,535-byte ceiling for the outer wrap's plaintext. See +`packages/core/src/transforms/chunk.ts` for the derivation. + +### Receiver + +`RelayClient` owns a `ChunkReassembler` instance. Chunks pass the same peer filter and +`lastProcessedTimestamp` dedup as any other message, then are routed to the reassembler. When all +chunks for a `msgId` have arrived (in any order), the bytes are concatenated, decoded, and the +resulting `ProtocolMessage` is handed to the application-level handler — indistinguishable to +the application from an unchunked message of equivalent size. + +The reassembler holds two maps with TTL eviction: + +- **in-flight buffers** keyed by `msgId` — collects chunks until complete; default TTL 120 s, + sized to comfortably fit a ~35-chunk 2 MB response under real relay latency. +- **completed** — tracks `msgId`s we've already delivered, for the same TTL, so late-arriving + duplicates (e.g. cross-subscription replay after reconnect) don't spawn a second reassembly and + double-deliver. + +A background sweeper runs every 10 s while connected and evicts expired entries. Reassembly state +is not persisted — reconnects rely on the relay replaying events to complete any in-flight +transfer (see below). + +### Failure modes + +| Failure | Behaviour | +|---|---| +| Dapp reloads mid-send | Dapp on reload has no in-flight state. User retries → fresh `msgId`, all chunks re-sent. Wallet's partial buffer for the old `msgId` expires via TTL. | +| Wallet reloads mid-receive | Subscription filter has no `since` clause, so on reconnect the relay re-delivers all events addressed to the wallet. Chunks reassemble fresh. Works as long as the chunks are still within the relay's retention window. | +| Network blip / all relays reject a chunk | `Promise.allSettled` on publish treats each chunk identically to any other single message. If all relays reject, `publishMessage` throws and `emitDisconnect` fires — same posture as today's sign-request failure. | +| Relay prunes mid-reassembly | Receiver's partial times out via TTL. Application sees no response; user retries — same failure mode the protocol already has for any lost sign request. | + +No persistence layer, no per-chunk ACKs, no new round trips. The assembled `ProtocolMessage`'s +own response (e.g. `sign_transaction_response`) is the effective end-to-end ACK. + +### Backward compatibility + +| Dapp | Wallet | Outcome | +|---|---|---| +| Old | Old | Unchanged. Oversized message fails at nostr-tools with the raw NIP-44 error. | +| **New** | Old | Dapp detects absent `chunk` capability and throws a clear upgrade-guidance error instead of the cryptic NIP-44 error. | +| Old | **New** | Wallet detects absent `chunk` capability and throws a clear error symmetrically. | +| **New** | **New** | Both advertise, both enable, oversized messages chunk-and-reassemble transparently. Application code is unchanged. | + +Upgrading the library on both sides is sufficient — no adapter-interface changes, no new +configuration, no capability opt-in. The extension is always-on when supported. diff --git a/package-lock.json b/package-lock.json index c097f3f..ca01763 100644 --- a/package-lock.json +++ b/package-lock.json @@ -3401,7 +3401,7 @@ }, "packages/core": { "name": "@wizardconnect/core", - "version": "0.1.2", + "version": "0.2.0", "dependencies": { "@bch-wc2/interfaces": "^0.0.8", "@bitauth/libauth": "^3.1.0-next.2", @@ -3418,7 +3418,7 @@ }, "packages/dapp": { "name": "@wizardconnect/dapp", - "version": "0.1.2", + "version": "0.2.0", "dependencies": { "@wizardconnect/core": "*", "eventemitter3": "^5.0.1" @@ -3430,7 +3430,7 @@ }, "packages/react": { "name": "@wizardconnect/react", - "version": "0.1.0", + "version": "0.2.0", "dependencies": { "@wizardconnect/core": "*", "@wizardconnect/dapp": "*", @@ -3502,7 +3502,7 @@ }, "packages/test-cli": { "name": "@wizardconnect/test-cli", - "version": "0.1.0", + "version": "0.2.0", "dependencies": { "@bitauth/libauth": "^3.1.0-next.2", "@wizardconnect/core": "*", @@ -3522,7 +3522,7 @@ }, "packages/wallet": { "name": "@wizardconnect/wallet", - "version": "0.1.2", + "version": "0.2.0", "dependencies": { "@bitauth/libauth": "^3.1.0-next.2", "@wizardconnect/core": "*", diff --git a/packages/core/package.json b/packages/core/package.json index f608080..7a7b5c4 100644 --- a/packages/core/package.json +++ b/packages/core/package.json @@ -1,6 +1,6 @@ { "name": "@wizardconnect/core", - "version": "0.1.2", + "version": "0.2.0", "type": "module", "description": "Transport and protocol primitives for WizardConnect", "repository": { diff --git a/packages/core/src/index.ts b/packages/core/src/index.ts index 0260ced..e7d6d1f 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -31,6 +31,12 @@ export type { } from "./key-exchange.js"; export * from "./protocols/hdwalletv1.js"; export * from "./protocols/base.js"; +export { + CHUNK_EXTENSION_NAME, + CHUNK_EXTENSION_VERSION, + chunkExtensionAdvertisement, + peerSupportsChunk, +} from "./transforms/chunk.js"; export { binToHex, hexToBin, diff --git a/packages/core/src/protocols/base.ts b/packages/core/src/protocols/base.ts index 70dcf45..5f1ac74 100644 --- a/packages/core/src/protocols/base.ts +++ b/packages/core/src/protocols/base.ts @@ -29,6 +29,10 @@ export interface DappReadyMessage extends ProtocolMessage { wallet_discovered: boolean; dapp_name?: string; dapp_icon?: string; + /// Transport-level extensions this dapp supports. Presence of a key = support. + /// Distinct from Hdwalletv1Session.extensions which is protocol-level. + /// See docs/transport.md. + extensions?: Record; } export interface WalletReadyMessage extends ProtocolMessage { @@ -48,6 +52,10 @@ export interface WalletReadyMessage extends ProtocolMessage { public_key: string; /// Echo of the shared secret from the connection URI (hex, 8 bytes). MITM prevention. secret: string; + /// Transport-level extensions this wallet supports. Presence of a key = support. + /// Distinct from Hdwalletv1Session.extensions which is protocol-level. + /// See docs/transport.md. + extensions?: Record; } export function isDappReadyMessage(msg: unknown): msg is DappReadyMessage { diff --git a/packages/core/src/protocols/hdwalletv1.ts b/packages/core/src/protocols/hdwalletv1.ts index 6752e73..b50d132 100644 --- a/packages/core/src/protocols/hdwalletv1.ts +++ b/packages/core/src/protocols/hdwalletv1.ts @@ -25,6 +25,10 @@ export enum RelayMsgAction { SignCancel = "sign_cancel", /// Courtesy notification: one side is closing the connection. Disconnect = "disconnect", + /// Transport-level: carries one slice of a message that exceeds NIP-44's + /// 65,535-byte plaintext ceiling. See ChunkMessage and docs/transport.md. + /// Not tied to hdwalletv1 semantics — applies to any application protocol. + Chunk = "chunk", } export interface ProtocolMessage { @@ -146,6 +150,20 @@ export interface SignCancelMessage extends ProtocolMessage { reason?: string; } +/// Transport-level message carrying one slice of a larger ProtocolMessage. +/// +/// All chunks of one logical message share the same msgId and time. +/// `data` is a base64-encoded slice of the UTF-8 bytes of +/// JSON.stringify(originalMessage); concatenate slices in `index` order, +/// base64-decode, UTF-8-decode, then JSON.parse to reconstruct. +export interface ChunkMessage extends ProtocolMessage { + action: RelayMsgAction.Chunk; + msgId: string; + index: number; + total: number; + data: string; +} + // Type guard functions export function isProtocolMessage(payload: any): payload is ProtocolMessage { @@ -193,3 +211,20 @@ export function isSignCancelMessage(msg: any): msg is SignCancelMessage { typeof msg.sequence === "number" ); } + +export function isChunkMessage(msg: any): msg is ChunkMessage { + return ( + msg && + typeof msg === "object" && + msg.action === RelayMsgAction.Chunk && + typeof msg.msgId === "string" && + typeof msg.index === "number" && + typeof msg.total === "number" && + typeof msg.data === "string" && + Number.isInteger(msg.index) && + Number.isInteger(msg.total) && + msg.total >= 1 && + msg.index >= 0 && + msg.index < msg.total + ); +} diff --git a/packages/core/src/relay-client.test.ts b/packages/core/src/relay-client.test.ts index df346c7..a5669e4 100644 --- a/packages/core/src/relay-client.test.ts +++ b/packages/core/src/relay-client.test.ts @@ -218,3 +218,88 @@ describe("RelayClient — isConnected", () => { expect(client.isConnected()).toBe(false); }); }); + +// --------------------------------------------------------------------------- +// Chunking — sender path +// --------------------------------------------------------------------------- + +describe("RelayClient — chunk sender path", () => { + afterEach(() => { + vi.restoreAllMocks(); + }); + + it("does not chunk small messages (regression guard)", async () => { + const { pool, triggerEose, publishMock } = makeMockPool(); + const client = makeClient(pool); + client.setPeerCapabilities({ chunk: true }); + await client.connect(); + triggerEose(); + + await client.relay({ + action: "dapp_ready" as any, + time: Math.floor(Date.now() / 1000), + }); + + expect(publishMock).toHaveBeenCalledTimes(1); + }); + + it("throws a clear error for oversized messages when peer lacks chunk support", async () => { + const { pool, triggerEose } = makeMockPool(); + const client = makeClient(pool); + // Intentionally NOT calling setPeerCapabilities({chunk: true}) + client.on("disconnect", () => {}); // prevent unhandled + await client.connect(); + triggerEose(); + + const huge = { + action: "sign_transaction_request" as any, + time: Math.floor(Date.now() / 1000), + payload: "x".repeat(200_000), + }; + + await expect(client.relay(huge)).rejects.toThrow( + /does not advertise the 'chunk' transport extension/, + ); + }); + + it("splits oversized messages into multiple publish calls when peer supports chunking", async () => { + const { pool, triggerEose, publishMock } = makeMockPool(); + const client = makeClient(pool); + client.setPeerCapabilities({ chunk: true }); + await client.connect(); + triggerEose(); + + const huge = { + action: "sign_transaction_request" as any, + time: Math.floor(Date.now() / 1000), + payload: "x".repeat(200_000), + }; + + await client.relay(huge); + + // 200 KB payload should produce multiple chunks; each chunk is one publish. + expect(publishMock.mock.calls.length).toBeGreaterThan(1); + }); + + it("setPeerCapabilities only updates provided keys", async () => { + const { pool, triggerEose } = makeMockPool(); + const client = makeClient(pool); + await client.connect(); + triggerEose(); + + // Enable, then call with empty object — should not disable + client.setPeerCapabilities({ chunk: true }); + client.setPeerCapabilities({}); + + const { publishMock } = makeMockPool(); // fresh count - not really needed + void publishMock; + + const huge = { + action: "sign_transaction_request" as any, + time: Math.floor(Date.now() / 1000), + payload: "x".repeat(200_000), + }; + // Should not throw, because chunk capability remains enabled + await expect(client.relay(huge)).resolves.toBeUndefined(); + }); +}); diff --git a/packages/core/src/relay-client.ts b/packages/core/src/relay-client.ts index 91967a7..d2c84a9 100644 --- a/packages/core/src/relay-client.ts +++ b/packages/core/src/relay-client.ts @@ -14,6 +14,8 @@ import WebSocket from "isomorphic-ws"; import { binToHex, hash256, secp256k1 } from "@bitauth/libauth"; import { EventEmitter } from "eventemitter3"; import { + ChunkMessage, + isChunkMessage, isProtocolMessage, ProtocolMessage, RelayMsgAction, @@ -21,6 +23,11 @@ import { import { deriveNostrPublicKey } from "./utilnostr.js"; import { MessageQueue } from "./message-queue.js"; import { debug, error as logError, Scope } from "./log.js"; +import { + ChunkReassembler, + needsChunking, + splitIntoChunks, +} from "./transforms/chunk.js"; useWebSocketImplementation(WebSocket); @@ -48,6 +55,14 @@ export class RelayClient extends EventEmitter { private disconnecting: boolean = false; + /// Capability flag: peer advertised support for the `chunk` transport + /// extension in its dapp_ready / wallet_ready. Set via setPeerCapabilities. + private peerSupportsChunk: boolean = false; + + /// Receiver-side reassembly buffer. Always active — if no chunks arrive it + /// stays empty. Started in connect(), stopped in disconnect(). + private reassembler: ChunkReassembler; + private sequence: number = Math.floor( Math.random() * (Number.MAX_SAFE_INTEGER - 500_000), ); @@ -81,6 +96,15 @@ export class RelayClient extends EventEmitter { logActivity: this.config.logNetworkActivity, }); + // Reassembled messages take the same post-decryption path as unchunked ones. + // Chunks themselves have already passed the peer filter and timestamp dedup + // on ingress (see routeIncoming), so the assembled message is handed directly + // to the application-level handler. + this.reassembler = new ChunkReassembler( + (msg) => this.handleRelayMessage(msg), + !!this.config.logNetworkActivity, + ); + this.myPubkey = unwrap( secp256k1.derivePublicKeyCompressed(this.config.signerPrivateKey), ); @@ -107,6 +131,16 @@ export class RelayClient extends EventEmitter { this.emit("paired"); } + /// Set transport-level capability flags based on the peer's advertisement in + /// its dapp_ready / wallet_ready `extensions` field. Called by the connection + /// manager after the handshake. New capability keys are additive — callers + /// may omit any they don't set. + setPeerCapabilities(caps: { chunk?: boolean }): void { + if (caps.chunk !== undefined) { + this.peerSupportsChunk = caps.chunk; + } + } + getPublicKey(): Uint8Array { return this.myPubkey; } @@ -134,6 +168,7 @@ export class RelayClient extends EventEmitter { } this.disconnecting = false; + this.reassembler.start(); if (this.lastProcessedTimestamp === 0) { this.lastProcessedTimestamp = Math.floor(Date.now() / 1000) - 2; @@ -194,6 +229,7 @@ export class RelayClient extends EventEmitter { async disconnect(): Promise { this.lastProcessedTimestamp = Math.floor(Date.now() / 1000); this.messageQueue.setNotReady(); + this.reassembler.stop(); if (this.readyTimeoutId) { clearTimeout(this.readyTimeoutId); @@ -233,12 +269,53 @@ export class RelayClient extends EventEmitter { } private async publishMessage(message: ProtocolMessage): Promise { - this.netlog("send", message.action); + const serialized = JSON.stringify(message); + + if (!needsChunking(serialized)) { + return this.publishSerialized(message.action, serialized); + } + + // Oversized. Chunk if the peer supports it; otherwise fail loudly with + // an actionable message (replacing nostr-tools' cryptic plaintext-size error). + if (!this.peerSupportsChunk) { + const err = new Error( + `Cannot send ${message.action}: message is larger than NIP-44's 65,535-byte ceiling ` + + `and the peer does not advertise the 'chunk' transport extension. ` + + `Please update the connected wallet/dapp to a version that supports chunked messages.`, + ); + if (this.config.logNetworkActivity) { + logError(Scope.Relay, err.message); + } + throw err; + } + + const chunks = splitIntoChunks(serialized); + if (this.config.logNetworkActivity) { + debug( + Scope.Relay, + `Chunking ${message.action}: ${chunks.length} chunks (serialized ~${serialized.length} bytes)`, + ); + } + for (const chunk of chunks) { + await this.publishSerialized( + `${message.action}[chunk ${chunk.index + 1}/${chunk.total}]`, + JSON.stringify(chunk), + ); + } + } + + /// Wrap and publish one gift-wrap event. Used for both unchunked messages + /// and individual chunks. `displayAction` is only used for logs. + private async publishSerialized( + displayAction: string, + serialized: string, + ): Promise { + this.netlog("send", displayAction); const wrapped = wrapEvent( { kind: KIND_PRIVATE_DIRECT_MESSAGE, - content: JSON.stringify(message), + content: serialized, created_at: Math.floor(Date.now() / 1000), tags: [["p", this.pairedPubkeyHex]], }, @@ -257,7 +334,7 @@ export class RelayClient extends EventEmitter { for (const r of rejected) { logError( Scope.Relay, - `Failed to publish ${message.action} to a relay:`, + `Failed to publish ${displayAction} to a relay:`, (r as PromiseRejectedResult).reason, ); } @@ -265,7 +342,7 @@ export class RelayClient extends EventEmitter { if (fulfilled.length === 0) { const error = new Error( - `Failed to publish ${message.action} to all relays`, + `Failed to publish ${displayAction} to all relays`, ); if (this.config.logNetworkActivity) { logError(Scope.Relay, error.message); @@ -277,7 +354,7 @@ export class RelayClient extends EventEmitter { if (this.config.logNetworkActivity) { debug( Scope.Relay, - `Published message ${message.action} to ${fulfilled.length}/${results.length} relay(s)`, + `Published message ${displayAction} to ${fulfilled.length}/${results.length} relay(s)`, ); } } @@ -286,68 +363,27 @@ export class RelayClient extends EventEmitter { try { const rumor = unwrapEvent(wrappedEvent, this.config.signerPrivateKey); - if (rumor.kind === KIND_PRIVATE_DIRECT_MESSAGE) { - let payload: ProtocolMessage; - try { - payload = JSON.parse(rumor.content); - } catch (e) { - if (this.config.logNetworkActivity) { - logError( - Scope.Relay, - "Failed to parse message content as JSON:", - e, - ); - } - return; - } - - if ( - !payload.time || - (this.lastProcessedTimestamp > 0 && - payload.time < this.lastProcessedTimestamp) - ) { - if (this.config.logNetworkActivity) { - debug( - Scope.Relay, - `Ignoring already-processed message (time: ${payload.time}, action: ${payload.action}, last processed: ${this.lastProcessedTimestamp})`, - ); - } - return; - } - - // wallet_ready carries the key exchange data (public_key + secret) so it - // must bypass the peer filter — the dapp doesn't know the wallet's pubkey yet. - const isKeyExchangeMessage = - payload.action === RelayMsgAction.WalletReady; - - if (!isKeyExchangeMessage && this.config.pairedPublicKey) { - const pairedNostrPubkey = - this.config.pairedPublicKey.length === 33 - ? binToHex(this.config.pairedPublicKey.slice(1)) - : binToHex(this.config.pairedPublicKey); - if (rumor.pubkey !== pairedNostrPubkey) { - if (this.config.logNetworkActivity) { - debug( - Scope.Relay, - `Ignoring '${payload.action}' message from unknown peer: ${rumor.pubkey} (expected: ${pairedNostrPubkey})`, - ); - } - return; - } - } - - if (this.config.logNetworkActivity) { - debug(Scope.Relay, `Received message ${payload.action} from relay`); - } - this.handleRelayMessage(payload); - } else { + if (rumor.kind !== KIND_PRIVATE_DIRECT_MESSAGE) { if (this.config.logNetworkActivity) { debug( Scope.Relay, `Ignoring non-PrivateDirectMessage, kind: ${rumor.kind}`, ); } + return; } + + let payload: ProtocolMessage; + try { + payload = JSON.parse(rumor.content); + } catch (e) { + if (this.config.logNetworkActivity) { + logError(Scope.Relay, "Failed to parse message content as JSON:", e); + } + return; + } + + this.routeIncoming(payload, rumor.pubkey); } catch (error) { if (this.config.logNetworkActivity) { logError(Scope.Relay, "Error handling incoming message:", error); @@ -356,6 +392,66 @@ export class RelayClient extends EventEmitter { } } + /// Apply timestamp dedup + peer filter, then dispatch to chunk reassembly + /// or the application-level handler. Called from handleWrappedEvent (one + /// path: unwrap → route). Kept separate to keep handleWrappedEvent focused + /// on decryption and to allow future transport-layer transforms to invoke + /// this path with already-decoded payloads. + private routeIncoming(payload: ProtocolMessage, fromPubkey: string): void { + if ( + !payload.time || + (this.lastProcessedTimestamp > 0 && + payload.time < this.lastProcessedTimestamp) + ) { + if (this.config.logNetworkActivity) { + debug( + Scope.Relay, + `Ignoring already-processed message (time: ${payload.time}, action: ${payload.action}, last processed: ${this.lastProcessedTimestamp})`, + ); + } + return; + } + + // wallet_ready carries the key exchange data (public_key + secret) so it + // must bypass the peer filter — the dapp doesn't know the wallet's pubkey yet. + const isKeyExchangeMessage = payload.action === RelayMsgAction.WalletReady; + + if (!isKeyExchangeMessage && this.config.pairedPublicKey) { + const pairedNostrPubkey = + this.config.pairedPublicKey.length === 33 + ? binToHex(this.config.pairedPublicKey.slice(1)) + : binToHex(this.config.pairedPublicKey); + if (fromPubkey !== pairedNostrPubkey) { + if (this.config.logNetworkActivity) { + debug( + Scope.Relay, + `Ignoring '${payload.action}' message from unknown peer: ${fromPubkey} (expected: ${pairedNostrPubkey})`, + ); + } + return; + } + } + + // Transport-level branch: a ChunkMessage is routed to the reassembler, + // which will emit a reassembled ProtocolMessage via handleRelayMessage + // once all pieces arrive. + if (isChunkMessage(payload)) { + if (this.config.logNetworkActivity) { + debug( + Scope.Relay, + `Received chunk ${payload.index + 1}/${payload.total} (msgId=${payload.msgId})`, + ); + } + this.reassembler.ingest(payload as ChunkMessage); + return; + } + + if (this.config.logNetworkActivity) { + debug(Scope.Relay, `Received message ${payload.action} from relay`); + } + this.handleRelayMessage(payload); + } + isConnected(): boolean { return this.subscription !== null; } diff --git a/packages/core/src/transforms/chunk.test.ts b/packages/core/src/transforms/chunk.test.ts new file mode 100644 index 0000000..060ded6 --- /dev/null +++ b/packages/core/src/transforms/chunk.test.ts @@ -0,0 +1,322 @@ +import { describe, it, expect, vi } from "vitest"; +import { + ChunkReassembler, + CHUNK_REQUIRED_BYTES, + CHUNK_RAW_BYTES, + REASSEMBLY_TTL_MS, + chunkExtensionAdvertisement, + needsChunking, + peerSupportsChunk, + splitIntoChunks, +} from "./chunk.js"; +import { + ChunkMessage, + ProtocolMessage, + RelayMsgAction, +} from "../protocols/hdwalletv1.js"; + +// --------------------------------------------------------------------------- +// Pure functions +// --------------------------------------------------------------------------- + +describe("chunk extension advertisement helpers", () => { + it("chunkExtensionAdvertisement returns a version-tagged object", () => { + const adv = chunkExtensionAdvertisement(); + expect(adv).toHaveProperty("version"); + expect(typeof adv.version).toBe("number"); + }); + + it("peerSupportsChunk true iff `chunk` key present", () => { + expect(peerSupportsChunk(undefined)).toBe(false); + expect(peerSupportsChunk({})).toBe(false); + expect(peerSupportsChunk({ chunk: { version: 1 } })).toBe(true); + expect(peerSupportsChunk({ chunk: {} })).toBe(true); + // presence of other keys doesn't imply chunk support + expect(peerSupportsChunk({ compress: {} })).toBe(false); + }); +}); + +describe("needsChunking", () => { + it("returns false for small payloads", () => { + expect(needsChunking("hello")).toBe(false); + expect(needsChunking("a".repeat(1000))).toBe(false); + }); + + it("returns true when UTF-8 bytes exceed the ceiling-minus-overhead", () => { + expect(needsChunking("a".repeat(CHUNK_REQUIRED_BYTES + 1))).toBe(true); + }); + + it("returns false right at the threshold", () => { + // ASCII: 1 byte per char. CHUNK_REQUIRED_BYTES chars fits. + expect(needsChunking("a".repeat(CHUNK_REQUIRED_BYTES))).toBe(false); + }); + + it("accounts for multibyte UTF-8 (emoji)", () => { + // "🦀" = 4 UTF-8 bytes. 20,000 crabs = 80,000 bytes > CHUNK_REQUIRED_BYTES. + const s = "🦀".repeat(20_000); + expect(needsChunking(s)).toBe(true); + }); +}); + +// --------------------------------------------------------------------------- +// Splitter +// --------------------------------------------------------------------------- + +describe("splitIntoChunks", () => { + it("produces one chunk for tiny input", () => { + const chunks = splitIntoChunks(JSON.stringify({ hello: "world" })); + expect(chunks).toHaveLength(1); + expect(chunks[0].index).toBe(0); + expect(chunks[0].total).toBe(1); + }); + + it("shares one msgId and one time across chunks", () => { + const big = "x".repeat(CHUNK_RAW_BYTES * 5); + const chunks = splitIntoChunks(big); + expect(chunks.length).toBeGreaterThan(1); + const firstId = chunks[0].msgId; + const firstTime = chunks[0].time; + for (const c of chunks) { + expect(c.msgId).toBe(firstId); + expect(c.time).toBe(firstTime); + } + }); + + it("every chunk passes the NIP-17 gift-wrap round-trip without exceeding NIP-44", async () => { + // The real constraint isn't just that the chunk's JSON fits in 65,535 + // plaintext bytes — it's that the OUTER gift-wrap's plaintext (which + // contains the seal, which contains the 4/3×-expanded encrypted rumor) + // also fits. This test exercises the full wrapEvent path. + const { wrapEvent } = await import("nostr-tools/nip59"); + const { generateSecretKey, getPublicKey } = + await import("nostr-tools/pure"); + const senderPriv = generateSecretKey(); + const recipPub = getPublicKey(generateSecretKey()); + + const big = "x".repeat(CHUNK_RAW_BYTES * 10); + const chunks = splitIntoChunks(big); + for (const c of chunks) { + expect(() => + wrapEvent( + { + kind: 14, + content: JSON.stringify(c), + created_at: Math.floor(Date.now() / 1000), + tags: [["p", recipPub]], + }, + senderPriv, + recipPub, + ), + ).not.toThrow(); + } + }); + + it("uses sequential indices starting at 0", () => { + const big = "x".repeat(CHUNK_RAW_BYTES * 4); + const chunks = splitIntoChunks(big); + chunks.forEach((c, i) => { + expect(c.index).toBe(i); + expect(c.total).toBe(chunks.length); + }); + }); + + it("accepts caller-supplied msgId and time", () => { + const chunks = splitIntoChunks("hello", { + msgId: "fixed-id", + time: 12345, + }); + expect(chunks[0].msgId).toBe("fixed-id"); + expect(chunks[0].time).toBe(12345); + }); +}); + +// --------------------------------------------------------------------------- +// Reassembler — round-trip +// --------------------------------------------------------------------------- + +describe("ChunkReassembler round-trip", () => { + it("reassembles a small message", () => { + const original: ProtocolMessage = { + action: "sign_transaction_request" as any, + time: 1000, + }; + const chunks = splitIntoChunks(JSON.stringify(original)); + const received: ProtocolMessage[] = []; + const r = new ChunkReassembler((m) => received.push(m), false); + for (const c of chunks) r.ingest(c); + expect(received).toHaveLength(1); + expect(received[0]).toEqual(original); + }); + + it("reassembles a ~2MB message (simulates signed tx hex response)", () => { + const original = { + action: "sign_transaction_response", + time: 1000, + sequence: 42, + signedTransaction: "ab".repeat(1_000_000), // 2 MB hex + }; + const chunks = splitIntoChunks(JSON.stringify(original)); + expect(chunks.length).toBeGreaterThan(30); + const received: any[] = []; + const r = new ChunkReassembler((m) => received.push(m), false); + for (const c of chunks) r.ingest(c); + expect(received).toHaveLength(1); + expect(received[0]).toEqual(original); + }); + + it("reassembles when chunks arrive in reverse order", () => { + const original = { + action: "sign_transaction_request", + time: 1, + payload: "x".repeat(CHUNK_RAW_BYTES * 3), + }; + const chunks = splitIntoChunks(JSON.stringify(original)); + const received: any[] = []; + const r = new ChunkReassembler((m) => received.push(m), false); + for (let i = chunks.length - 1; i >= 0; i--) r.ingest(chunks[i]); + expect(received).toHaveLength(1); + expect(received[0]).toEqual(original); + }); + + it("reassembles when chunks arrive in shuffled order", () => { + const original = { + action: "sign_transaction_request", + time: 1, + payload: "x".repeat(CHUNK_RAW_BYTES * 4), + }; + const chunks = splitIntoChunks(JSON.stringify(original)); + const shuffled = [...chunks].sort(() => Math.random() - 0.5); + const received: any[] = []; + const r = new ChunkReassembler((m) => received.push(m), false); + for (const c of shuffled) r.ingest(c); + expect(received).toHaveLength(1); + expect(received[0]).toEqual(original); + }); + + it("round-trips multibyte UTF-8 content", () => { + const original = { + action: "custom_action", + time: 1, + crab: "🦀".repeat(20_000), // 80 KB of 4-byte codepoints + japanese: "こんにちは世界".repeat(5_000), + }; + const chunks = splitIntoChunks(JSON.stringify(original)); + const received: any[] = []; + const r = new ChunkReassembler((m) => received.push(m), false); + for (const c of chunks) r.ingest(c); + expect(received).toHaveLength(1); + expect(received[0]).toEqual(original); + }); +}); + +// --------------------------------------------------------------------------- +// Reassembler — duplicate / malformed / TTL +// --------------------------------------------------------------------------- + +describe("ChunkReassembler edge cases", () => { + it("duplicate chunks are idempotent (delivers exactly once)", () => { + const original = { + action: "a", + time: 1, + x: "y".repeat(CHUNK_RAW_BYTES * 2), + }; + const chunks = splitIntoChunks(JSON.stringify(original)); + const received: any[] = []; + const r = new ChunkReassembler((m) => received.push(m), false); + for (const c of chunks) r.ingest(c); + // Replay everything a second time + for (const c of chunks) r.ingest(c); + expect(received).toHaveLength(1); + }); + + it("does not deliver if a chunk is missing", () => { + const original = { + action: "a", + time: 1, + x: "y".repeat(CHUNK_RAW_BYTES * 3), + }; + const chunks = splitIntoChunks(JSON.stringify(original)); + const received: any[] = []; + const r = new ChunkReassembler((m) => received.push(m), false); + for (let i = 0; i < chunks.length - 1; i++) r.ingest(chunks[i]); + expect(received).toHaveLength(0); + }); + + it("drops a chunk whose total disagrees with the in-flight entry", () => { + const chunkA: ChunkMessage = { + action: RelayMsgAction.Chunk, + time: 1, + msgId: "m", + index: 0, + total: 3, + data: "AAAA", + }; + const chunkBadTotal: ChunkMessage = { ...chunkA, index: 1, total: 99 }; + const received: any[] = []; + const r = new ChunkReassembler((m) => received.push(m), false); + r.ingest(chunkA); + r.ingest(chunkBadTotal); + expect(r.bufferCount).toBe(1); + expect(received).toHaveLength(0); + }); + + it("reassembled payload that is not a valid ProtocolMessage is dropped", () => { + // Craft a single chunk carrying junk JSON + const junk = JSON.stringify({ not: "a protocol message" }); + const chunks = splitIntoChunks(junk); + const received: any[] = []; + const r = new ChunkReassembler((m) => received.push(m), false); + for (const c of chunks) r.ingest(c); + expect(received).toHaveLength(0); + }); + + it("reassembled payload with invalid base64/UTF-8 is dropped without crashing", () => { + const chunk: ChunkMessage = { + action: RelayMsgAction.Chunk, + time: 1, + msgId: "bad", + index: 0, + total: 1, + data: "!!!not-base64!!!", + }; + const received: any[] = []; + const r = new ChunkReassembler((m) => received.push(m), false); + expect(() => r.ingest(chunk)).not.toThrow(); + expect(received).toHaveLength(0); + }); + + it("TTL sweeper evicts incomplete entries", () => { + let clock = 0; + const r = new ChunkReassembler( + () => {}, + false, + () => clock, + ); + const original = { + action: "a", + time: 1, + x: "y".repeat(CHUNK_RAW_BYTES * 3), + }; + const chunks = splitIntoChunks(JSON.stringify(original)); + r.ingest(chunks[0]); + expect(r.bufferCount).toBe(1); + + // Advance clock past TTL + clock += REASSEMBLY_TTL_MS + 1; + r.sweep(); + expect(r.bufferCount).toBe(0); + }); + + it("start() installs a periodic sweeper that stop() clears", () => { + vi.useFakeTimers(); + const r = new ChunkReassembler(() => {}, false); + r.start(); + // Shouldn't throw on repeat start + r.start(); + vi.advanceTimersByTime(100_000); + r.stop(); + // Shouldn't throw on repeat stop + r.stop(); + vi.useRealTimers(); + }); +}); diff --git a/packages/core/src/transforms/chunk.ts b/packages/core/src/transforms/chunk.ts new file mode 100644 index 0000000..5688295 --- /dev/null +++ b/packages/core/src/transforms/chunk.ts @@ -0,0 +1,305 @@ +// Copyright (C) 2026 Whiterun LLC, +// This software is licensed under the GNU Lesser General Public License (LGPL), version 3.0 or later. +// A copy of the license can be found in the LICENSE file or at https://www.gnu.org/licenses/lgpl-3.0.html + +/** + * Transport-level chunking. + * + * NIP-44 caps plaintext at 65,535 bytes (the plaintext length is encoded as a + * U16BE prefix in the wire format — structural, not a configurable guardrail). + * This module splits oversized ProtocolMessages into a sequence of + * ChunkMessages that each fit under the ceiling, and reassembles them on the + * receiver. + * + * Fire-and-forget: chunks are published individually through the same + * multi-relay path as any other message, with no per-chunk ACKs. Receiver + * buffers by msgId, applies a TTL, and delivers the assembled ProtocolMessage + * once complete. + * + * See docs/transport.md for wire format and failure modes. + */ + +import { + ChunkMessage, + ProtocolMessage, + RelayMsgAction, + isProtocolMessage, +} from "../protocols/hdwalletv1.js"; +import { debug, error as logError, Scope } from "../log.js"; + +/// Extension key advertised in base-level `extensions` on dapp_ready / wallet_ready. +export const CHUNK_EXTENSION_NAME = "chunk"; +export const CHUNK_EXTENSION_VERSION = 1; + +/// NIP-44 plaintext hard limit. +export const NIP44_MAX_PLAINTEXT = 65535; + +/// NIP-17 gift-wrap applies NIP-44 encryption twice: once to the rumor +/// (inner) and once to the seal (outer). The 65,535-byte cap applies to the +/// plaintext of each layer. What we hand to wrapEvent becomes the rumor +/// content; that rumor is then JSON-serialized, padded (up to multiples of +/// 8192 for sizes in this range), encrypted, base64'd, and JSON-wrapped in +/// a seal — the outer wrap then encrypts that seal JSON, which must also +/// fit under 65,535 plaintext bytes. +/// +/// For a rumor content of C bytes (ASCII) the outer plaintext is roughly: +/// (event_shell≈200) + ceil(4/3 × (32 + pad(C + 200) + 32)) + seal_shell≈300 +/// With pad(~40 KB) = 40,960 this totals ≈55 KB — safely under the cap. +/// Above ~40,960 bytes of content the next pad multiple jumps to 49,152, +/// pushing the outer plaintext past the cap. + +/// Raw content bytes we'll pack into a single chunk's `data` field (before +/// base64). Sized so that after base64 expansion, envelope overhead, and the +/// two-layer gift-wrap expansion, the outer plaintext stays well under +/// NIP-44's 65,535-byte ceiling. See derivation above. +export const CHUNK_RAW_BYTES = 30000; + +/// When JSON.stringify(message)'s UTF-8 byte length exceeds this, chunking +/// is required. Derived from the same outer-wrap size analysis: anything +/// above ~40,760 bytes of content pushes the outer wrap plaintext over +/// 65,535. Conservative headroom built in. +export const CHUNK_REQUIRED_BYTES = 40000; + +/// How long an incomplete reassembly buffer survives without progress. +/// Sized for a ~35-chunk transfer (2 MB tx-hex response) under congested +/// relay conditions. +export const REASSEMBLY_TTL_MS = 120_000; + +/// How often the TTL sweeper runs. +export const SWEEP_INTERVAL_MS = 10_000; + +/// Returns the extension advertisement value for the `chunk` key in the +/// base-level `extensions` field of dapp_ready / wallet_ready. +export function chunkExtensionAdvertisement(): { version: number } { + return { version: CHUNK_EXTENSION_VERSION }; +} + +/// True iff the peer's `extensions` object advertises support for chunking. +export function peerSupportsChunk( + extensions: Record | undefined, +): boolean { + return !!extensions && extensions[CHUNK_EXTENSION_NAME] !== undefined; +} + +/// Measure UTF-8 byte length without allocating the full encoded buffer +/// (a shallow optimisation for very large payloads). +export function utf8ByteLength(s: string): number { + return new TextEncoder().encode(s).length; +} + +/// True iff a message of this serialized size requires chunking. +export function needsChunking(serialized: string): boolean { + return utf8ByteLength(serialized) > CHUNK_REQUIRED_BYTES; +} + +/// Split a serialized ProtocolMessage into ChunkMessages. All chunks share +/// one msgId and one `time` so reassembly is deterministic and the +/// reassembled message keeps a coherent timestamp for existing dedup logic. +export function splitIntoChunks( + serialized: string, + opts?: { msgId?: string; time?: number }, +): ChunkMessage[] { + const msgId = opts?.msgId ?? newMsgId(); + const time = opts?.time ?? Math.floor(Date.now() / 1000); + const utf8 = new TextEncoder().encode(serialized); + const b64 = bytesToBase64(utf8); + // base64 chars per chunk equivalent to CHUNK_RAW_BYTES raw bytes + const sliceChars = Math.ceil((CHUNK_RAW_BYTES * 4) / 3); + const total = Math.max(1, Math.ceil(b64.length / sliceChars)); + const chunks: ChunkMessage[] = []; + for (let i = 0; i < total; i++) { + chunks.push({ + action: RelayMsgAction.Chunk, + time, + msgId, + index: i, + total, + data: b64.slice(i * sliceChars, (i + 1) * sliceChars), + }); + } + return chunks; +} + +interface ReassemblyEntry { + total: number; + chunks: (string | undefined)[]; + received: number; + expiresAt: number; +} + +/** + * Receiver-side chunk reassembler. + * + * Owns a msgId-keyed buffer of partial messages, TTL-evicted by a periodic + * sweeper. When all chunks for a msgId are present the assembled + * ProtocolMessage is handed to `onComplete`. + * + * Not concurrency-safe — a single instance is owned by a single RelayClient. + */ +export class ChunkReassembler { + private buffers = new Map(); + /// msgIds that have already been delivered, kept for a short grace period + /// so late-arriving duplicate chunks (e.g. cross-subscription replay after + /// reconnect) don't spawn a second reassembly and double-deliver. + private completed = new Map(); + private sweeperId: ReturnType | null = null; + + constructor( + private readonly onComplete: (msg: ProtocolMessage) => void, + private readonly logActivity: boolean = true, + private readonly now: () => number = () => Date.now(), + ) {} + + start(): void { + if (this.sweeperId !== null) return; + this.sweeperId = setInterval(() => this.sweep(), SWEEP_INTERVAL_MS); + } + + stop(): void { + if (this.sweeperId !== null) { + clearInterval(this.sweeperId); + this.sweeperId = null; + } + this.buffers.clear(); + this.completed.clear(); + } + + /// For tests. Otherwise start() schedules this automatically. + sweep(): void { + const now = this.now(); + for (const [id, entry] of this.buffers) { + if (entry.expiresAt <= now) { + this.buffers.delete(id); + if (this.logActivity) { + debug( + Scope.Relay, + `Chunk reassembly timeout: ${id} (${entry.received}/${entry.total} received)`, + ); + } + } + } + for (const [id, expiresAt] of this.completed) { + if (expiresAt <= now) this.completed.delete(id); + } + } + + /// Ingest one chunk. If this completes the message, onComplete fires. + /// Duplicate chunks (same msgId/index) are idempotent. Malformed chunks + /// are dropped silently. + ingest(chunk: ChunkMessage): void { + if (this.completed.has(chunk.msgId)) { + // Already delivered this message; drop late duplicates. + return; + } + let entry = this.buffers.get(chunk.msgId); + if (!entry) { + entry = { + total: chunk.total, + chunks: new Array(chunk.total), + received: 0, + expiresAt: this.now() + REASSEMBLY_TTL_MS, + }; + this.buffers.set(chunk.msgId, entry); + } else if (entry.total !== chunk.total) { + // Protocol invariant violation — peer changed total mid-stream. Drop. + if (this.logActivity) { + logError( + Scope.Relay, + `Chunk total mismatch for ${chunk.msgId}: ${chunk.total} vs expected ${entry.total}`, + ); + } + return; + } + + if (entry.chunks[chunk.index] !== undefined) { + // Duplicate — already have this slot. SimplePool dedupes by event id + // within a subscription; this guards against the cross-subscription + // replay case after a reconnect. + return; + } + entry.chunks[chunk.index] = chunk.data; + entry.received++; + + if (entry.received !== entry.total) return; + + // Assemble and deliver. + this.buffers.delete(chunk.msgId); + // Grace period matches reassembly TTL — sufficient to catch any laggard + // duplicates from a slow relay that were in flight when we completed. + this.completed.set(chunk.msgId, this.now() + REASSEMBLY_TTL_MS); + let reassembled: ProtocolMessage; + try { + const fullB64 = entry.chunks.join(""); + const bytes = base64ToBytes(fullB64); + const json = new TextDecoder().decode(bytes); + const parsed = JSON.parse(json); + if (!isProtocolMessage(parsed)) { + if (this.logActivity) { + logError( + Scope.Relay, + `Reassembled chunk msgId=${chunk.msgId} is not a valid ProtocolMessage`, + ); + } + return; + } + reassembled = parsed; + } catch (e) { + if (this.logActivity) { + logError( + Scope.Relay, + `Chunk reassembly failed for msgId=${chunk.msgId}:`, + e, + ); + } + return; + } + + if (this.logActivity) { + debug( + Scope.Relay, + `Reassembled chunked message: action=${reassembled.action} chunks=${entry.total}`, + ); + } + this.onComplete(reassembled); + } + + /// Test helper: number of in-flight partial messages. + get bufferCount(): number { + return this.buffers.size; + } +} + +// --------------------------------------------------------------------------- +// Cross-platform base64 (no dependency on Buffer or DOM-specific APIs). +// --------------------------------------------------------------------------- + +function bytesToBase64(bytes: Uint8Array): string { + const CHUNK = 0x8000; + let binary = ""; + for (let i = 0; i < bytes.length; i += CHUNK) { + const slice = bytes.subarray(i, i + CHUNK); + binary += String.fromCharCode.apply(null, slice as unknown as number[]); + } + return btoa(binary); +} + +function base64ToBytes(b64: string): Uint8Array { + const binary = atob(b64); + const out = new Uint8Array(binary.length); + for (let i = 0; i < binary.length; i++) { + out[i] = binary.charCodeAt(i); + } + return out; +} + +function newMsgId(): string { + const g = globalThis as { crypto?: { randomUUID?: () => string } }; + if (g.crypto?.randomUUID) return g.crypto.randomUUID(); + const bytes = new Uint8Array(16); + const cryptoObj = ( + globalThis as { crypto?: { getRandomValues?: (b: Uint8Array) => void } } + ).crypto; + if (cryptoObj?.getRandomValues) cryptoObj.getRandomValues(bytes); + else for (let i = 0; i < 16; i++) bytes[i] = Math.floor(Math.random() * 256); + return Array.from(bytes, (b) => b.toString(16).padStart(2, "0")).join(""); +} diff --git a/packages/dapp/package.json b/packages/dapp/package.json index f3c5db3..75d11cd 100644 --- a/packages/dapp/package.json +++ b/packages/dapp/package.json @@ -1,6 +1,6 @@ { "name": "@wizardconnect/dapp", - "version": "0.1.2", + "version": "0.2.0", "type": "module", "description": "Dapp-side integration helpers for WizardConnect", "repository": { diff --git a/packages/dapp/src/dapp-connection-manager.test.ts b/packages/dapp/src/dapp-connection-manager.test.ts index 8a9ad4f..8ae3f1d 100644 --- a/packages/dapp/src/dapp-connection-manager.test.ts +++ b/packages/dapp/src/dapp-connection-manager.test.ts @@ -148,6 +148,7 @@ describe("DappConnectionManager", () => { relayed.push(msg); }), isKeyExchangeComplete: () => true, + setPeerCapabilities: vi.fn(), nextSequence: (() => { let seq = 0; return () => (seq += 2); diff --git a/packages/dapp/src/dapp-connection-manager.ts b/packages/dapp/src/dapp-connection-manager.ts index 1b3495c..a40ff55 100644 --- a/packages/dapp/src/dapp-connection-manager.ts +++ b/packages/dapp/src/dapp-connection-manager.ts @@ -21,6 +21,8 @@ import { childIndexOfPathName, isHdwalletv1Session, binToHex, + chunkExtensionAdvertisement, + peerSupportsChunk, } from "@wizardconnect/core"; import type { PathXpub, DappRelayResult } from "@wizardconnect/core"; import { DappPubkeyStateManager } from "./pubkey-state-manager.js"; @@ -409,6 +411,10 @@ export class DappConnectionManager extends EventEmitter { + let conn: ConnectionHandles; + const pending: PendingSignRequest[] = []; + + beforeAll(async () => { + conn = await setupConnection(TEST_RELAY_URL); + conn.wallet.manager.on("pendingSignRequest", (req) => pending.push(req)); + }, 30000); + + afterAll(() => { + conn?.cleanup(); + }); + + it("wallet_ready advertises the chunk transport extension", () => { + const wr = conn.dapp.walletReadyMessages[0]; + expect(wr.extensions?.chunk).toBeDefined(); + }); + + it("reassembles an 80 KB sign_transaction_request", async () => { + const bigHex = "ab".repeat(40_000); // 80 KB hex string + const msg: SignTransactionRequest = { + action: RelayMsgAction.SignTransactionRequest, + sequence: 100, + transaction: { + transaction: { inputs: [], outputs: [], version: 2, locktime: 0 }, + sourceOutputs: [], + userPrompt: bigHex, + broadcast: false, + }, + inputPaths: [], + time: Math.floor(Date.now() / 1000), + }; + await conn.dapp.client()!.relay(msg); + + await waitFor(() => pending.some((p) => p.request.sequence === 100), { + timeoutMs: 30000, + what: "80 KB sign request reassembled on wallet", + }); + + const got = pending.find((p) => p.request.sequence === 100)!; + expect(got.request.transaction.userPrompt).toBe(bigHex); + }, 60000); + + it("reassembles a 200 KB sign_transaction_request", async () => { + const bigHex = "cd".repeat(100_000); // 200 KB hex + const msg: SignTransactionRequest = { + action: RelayMsgAction.SignTransactionRequest, + sequence: 101, + transaction: { + transaction: { inputs: [], outputs: [], version: 2, locktime: 0 }, + sourceOutputs: [], + userPrompt: bigHex, + broadcast: false, + }, + inputPaths: [], + time: Math.floor(Date.now() / 1000), + }; + await conn.dapp.client()!.relay(msg); + + await waitFor(() => pending.some((p) => p.request.sequence === 101), { + timeoutMs: 45000, + what: "200 KB sign request reassembled on wallet", + }); + + const got = pending.find((p) => p.request.sequence === 101)!; + expect(got.request.transaction.userPrompt).toBe(bigHex); + expect(got.request.transaction.userPrompt!.length).toBe(200_000); + }, 60000); +}); + +describe("chunk extension — wallet → dapp oversized responses", () => { + let conn: ConnectionHandles; + + beforeAll(async () => { + conn = await setupConnection(TEST_RELAY_URL); + }, 30000); + + afterAll(() => { + conn?.cleanup(); + }); + + async function walletSend(msg: SignTransactionResponse): Promise { + // Reach through WalletConnectionManager internals to get the RelayClient + // directly and publish. A production app would return the response + // through its normal sign-approval flow; integration tests bypass that. + const m = conn.wallet.manager as unknown as { + connections: Map< + string, + { client: { relay: (m: unknown) => Promise } } + >; + }; + const connEntry = m.connections.get(conn.wallet.connectionId); + if (!connEntry?.client) throw new Error("wallet relay client not ready"); + await connEntry.client.relay(msg); + } + + it("reassembles a ~100 KB signed tx hex response", async () => { + const hex = "ef".repeat(50_000); // 100 KB hex + const msg: SignTransactionResponse = { + action: RelayMsgAction.SignTransactionResponse, + sequence: 200, + signedTransaction: hex, + time: Math.floor(Date.now() / 1000), + }; + await walletSend(msg); + + await waitFor( + () => + conn.dapp.messages.some( + (m) => + m.action === RelayMsgAction.SignTransactionResponse && + (m as SignTransactionResponse).sequence === 200, + ), + { + timeoutMs: 45000, + what: "100 KB sign response reassembled on dapp", + }, + ); + + const got = conn.dapp.messages.find( + (m) => + m.action === RelayMsgAction.SignTransactionResponse && + (m as SignTransactionResponse).sequence === 200, + ) as SignTransactionResponse; + expect(got.signedTransaction).toBe(hex); + }, 60000); + + it("reassembles a ~2 MB signed tx hex response (consensus-max case)", async () => { + const hex = "f0".repeat(1_000_000); // 2 MB hex + const msg: SignTransactionResponse = { + action: RelayMsgAction.SignTransactionResponse, + sequence: 201, + signedTransaction: hex, + time: Math.floor(Date.now() / 1000), + }; + await walletSend(msg); + + await waitFor( + () => + conn.dapp.messages.some( + (m) => + m.action === RelayMsgAction.SignTransactionResponse && + (m as SignTransactionResponse).sequence === 201, + ), + { + timeoutMs: 120000, + what: "2 MB sign response reassembled on dapp", + }, + ); + + const got = conn.dapp.messages.find( + (m) => + m.action === RelayMsgAction.SignTransactionResponse && + (m as SignTransactionResponse).sequence === 201, + ) as SignTransactionResponse; + expect(got.signedTransaction.length).toBe(2_000_000); + expect(got.signedTransaction).toBe(hex); + }, 180000); +}); diff --git a/packages/wallet/src/integration/helpers.ts b/packages/wallet/src/integration/helpers.ts index 86dd717..d3b0607 100644 --- a/packages/wallet/src/integration/helpers.ts +++ b/packages/wallet/src/integration/helpers.ts @@ -22,6 +22,8 @@ import { initiateDappRelay, RelayMsgAction, PROTOCOL_NAME, + chunkExtensionAdvertisement, + peerSupportsChunk, type RelayClient, type RelayUpdatePayload, type DappReadyMessage, @@ -110,6 +112,10 @@ export interface DappHandle { cleanup: () => void; /** All wallet_ready messages received. */ walletReadyMessages: WalletReadyMessage[]; + /** RelayClient for direct message publishing (populated after key exchange). */ + client: () => RelayClient | null; + /** All non-handshake messages received by the dapp. */ + messages: ProtocolMessage[]; } // ---- WalletHandle ----------------------------------------------------------- @@ -147,6 +153,7 @@ export async function setupConnection( // ---- Dapp side ---- const walletReadyMessages: WalletReadyMessage[] = []; + const allMessages: ProtocolMessage[] = []; let dappClient: RelayClient | null = null; let keyExchanged = false; @@ -163,6 +170,9 @@ export async function setupConnection( supported_protocols: [PROTOCOL_NAME], wallet_discovered: wd, time: Math.floor(Date.now() / 1000), + // Advertise transport-level chunking so the wallet can return + // oversized sign_transaction_response messages. + extensions: { chunk: chunkExtensionAdvertisement() }, }; await dappClient!.relay(msg); } @@ -182,9 +192,16 @@ export async function setupConnection( if (message.action === RelayMsgAction.WalletReady) { const msg = message as WalletReadyMessage; walletReadyMessages.push(msg); + // Mirror DappConnectionManager behavior: enable chunked outbound if + // the wallet advertises support. + dappClient!.setPeerCapabilities({ + chunk: peerSupportsChunk(msg.extensions), + }); if (!msg.dapp_discovered) { sendDappReady(true).catch(() => {}); } + } else { + allMessages.push(message); } }); @@ -217,6 +234,8 @@ export async function setupConnection( uri: dappRelay.uri, cleanup: dappRelay.cleanup, walletReadyMessages, + client: () => dappClient, + messages: allMessages, }; const walletHandle: WalletHandle = { manager, connectionId, adapter }; diff --git a/packages/wallet/src/wallet-connection-manager.ts b/packages/wallet/src/wallet-connection-manager.ts index a437980..a2e2bb5 100644 --- a/packages/wallet/src/wallet-connection-manager.ts +++ b/packages/wallet/src/wallet-connection-manager.ts @@ -22,6 +22,8 @@ import { Hdwalletv1Session, PROTOCOL_NAME, binToHex, + chunkExtensionAdvertisement, + peerSupportsChunk, } from "@wizardconnect/core"; import { WalletAdapter } from "./wallet-adapter.js"; import { DerivationPath } from "./derivation-path.js"; @@ -385,6 +387,15 @@ export class WalletConnectionManager extends EventEmitter