From fe0045be21a07198c157c72bc8bf847bd2c52dbb Mon Sep 17 00:00:00 2001 From: Benny Date: Thu, 30 Jul 2026 15:01:30 +0200 Subject: [PATCH 01/13] Bound MSB apply waits --- package-lock.json | 5 +-- package.json | 1 + src/config/config.js | 20 ++++++++++++ src/config/env.js | 6 ++++ src/index.js | 5 ++- src/msbClient.js | 28 ++++++++++++++-- src/operations/tx/index.js | 5 +++ tests/unit/applyGuards.test.js | 59 ++++++++++++++++++++++++++++++++++ 8 files changed, 124 insertions(+), 5 deletions(-) diff --git a/package-lock.json b/package-lock.json index e14f246..bcd0fad 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1,12 +1,12 @@ { "name": "trac-peer", - "version": "0.4.5", + "version": "0.4.6", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "trac-peer", - "version": "0.4.5", + "version": "0.4.6", "dependencies": { "@tracsystems/blake3": "^0.0.15", "assert": "npm:bare-node-assert", @@ -66,6 +66,7 @@ "process": "npm:bare-node-process", "protomux": "^3.10.1", "protomux-wakeup": "^2.4.0", + "rache": "1.0.0", "readline": "npm:bare-node-readline", "ready-resource": "1.1.2", "repl": "npm:bare-node-repl", diff --git a/package.json b/package.json index f9f7b5d..523372c 100644 --- a/package.json +++ b/package.json @@ -91,6 +91,7 @@ "process": "npm:bare-node-process", "protomux": "^3.10.1", "protomux-wakeup": "^2.4.0", + "rache": "1.0.0", "readline": "npm:bare-node-readline", "ready-resource": "1.1.2", "repl": "npm:bare-node-repl", diff --git a/src/config/config.js b/src/config/config.js index c8121e4..465af39 100644 --- a/src/config/config.js +++ b/src/config/config.js @@ -1,6 +1,9 @@ import b4a from "b4a"; import path from "path"; +const DEFAULT_HYPERBEE_CACHE_MAX_ENTRIES = 65_536; +const DEFAULT_MAX_MSB_SIGNED_LENGTH_FUTURE_DELTA = 1_000_000; + export class Config { #options; @@ -41,12 +44,29 @@ export class Config { } this.maxTxDelay = maxTxDelay; + const hyperbeeCacheMaxEntries = + this.#select("hyperbeeCacheMaxEntries", options, defaults) ?? + DEFAULT_HYPERBEE_CACHE_MAX_ENTRIES; + if (!Number.isSafeInteger(hyperbeeCacheMaxEntries) || hyperbeeCacheMaxEntries < 1) { + throw new Error("Peer: hyperbeeCacheMaxEntries must be a positive safe integer."); + } + this.hyperbeeCacheMaxEntries = hyperbeeCacheMaxEntries; + const maxMsbSignedLength = this.#select("maxMsbSignedLength", options, defaults); if (!Number.isSafeInteger(maxMsbSignedLength)) { throw new Error("Peer: maxMsbSignedLength must be a safe integer."); } this.maxMsbSignedLength = maxMsbSignedLength; + const maxMsbSignedLengthFutureDelta = + this.#select("maxMsbSignedLengthFutureDelta", options, defaults) ?? + DEFAULT_MAX_MSB_SIGNED_LENGTH_FUTURE_DELTA; + if (!Number.isSafeInteger(maxMsbSignedLengthFutureDelta) || + maxMsbSignedLengthFutureDelta < 0) { + throw new Error("Peer: maxMsbSignedLengthFutureDelta must be a non-negative safe integer."); + } + this.maxMsbSignedLengthFutureDelta = maxMsbSignedLengthFutureDelta; + const maxMsbApplyOperationBytes = this.#select("maxMsbApplyOperationBytes", options, defaults); if (!Number.isSafeInteger(maxMsbApplyOperationBytes)) { throw new Error("Peer: maxMsbApplyOperationBytes must be a safe integer."); diff --git a/src/config/env.js b/src/config/env.js index 4dc7b17..c0ee795 100644 --- a/src/config/env.js +++ b/src/config/env.js @@ -14,7 +14,9 @@ const configData = { storeName: "mainnet", txPoolMaxSize: 1_000, maxTxDelay: 60, + hyperbeeCacheMaxEntries: 65_536, maxMsbSignedLength: 1_000_000_000, + maxMsbSignedLengthFutureDelta: 1_000_000, maxMsbApplyOperationBytes: 1024 * 1024, enableInteractiveMode: true, enableBackgroundTasks: true, @@ -38,7 +40,9 @@ const configData = { storeName: "testnet", txPoolMaxSize: 1_000, maxTxDelay: 60, + hyperbeeCacheMaxEntries: 65_536, maxMsbSignedLength: 1_000_000_000, + maxMsbSignedLengthFutureDelta: 1_000_000, maxMsbApplyOperationBytes: 1024 * 1024, enableInteractiveMode: true, enableBackgroundTasks: true, @@ -62,7 +66,9 @@ const configData = { storeName: "peer", txPoolMaxSize: 1_000, maxTxDelay: 60, + hyperbeeCacheMaxEntries: 65_536, maxMsbSignedLength: 1_000_000_000, + maxMsbSignedLengthFutureDelta: 1_000_000, maxMsbApplyOperationBytes: 1024 * 1024, enableInteractiveMode: true, enableBackgroundTasks: false, diff --git a/src/index.js b/src/index.js index 79d6fd6..0114489 100644 --- a/src/index.js +++ b/src/index.js @@ -5,6 +5,7 @@ import ReadyResource from 'ready-resource'; import b4a from 'b4a'; import Hyperbee from 'hyperbee'; import Corestore from 'corestore'; +import Rache from 'rache'; import w from 'protomux-wakeup'; const wakeup = new w(); import Protomux from 'protomux' @@ -32,7 +33,9 @@ export class Peer extends ReadyResource { this.config = config; this.keyPair = null; - this.store = new Corestore(this.config.fullStoresDirectory); + this.store = new Corestore(this.config.fullStoresDirectory, { + globalCache: new Rache({ maxSize: this.config.hyperbeeCacheMaxEntries }) + }); this.msbClient = new MsbClient(msb); this.swarm = null; this.base = null; diff --git a/src/msbClient.js b/src/msbClient.js index 9c1452d..55e2279 100644 --- a/src/msbClient.js +++ b/src/msbClient.js @@ -97,11 +97,35 @@ export class MsbClient extends ReadyResource { return await this.#msb.network.tryConnect(pubKeyHex, role); } - async waitForSignedLengthAtLeast(targetSignedLength) { + async waitForSignedLengthAtLeast(targetSignedLength, { pollMs = 1_000 } = {}) { const core = this.#msb.state?.base?.view?.core ?? null; if (!core) throw new Error('MSB view core not available.'); + if (!Number.isSafeInteger(targetSignedLength) || targetSignedLength < 0) { + throw new Error('Invalid MSB signed length target.'); + } + if (!Number.isSafeInteger(pollMs) || pollMs < 1) { + throw new Error('Invalid MSB signed length wait poll interval.'); + } while (core.signedLength < targetSignedLength) { - await new Promise((resolve) => core.once('append', resolve)); + await new Promise((resolve) => { + const onAppend = () => { + cleanup(); + resolve(); + }; + const cleanup = () => { + clearTimeout(timer); + if (typeof core.off === 'function') { + core.off('append', onAppend); + } else if (typeof core.removeListener === 'function') { + core.removeListener('append', onAppend); + } + }; + const timer = setTimeout(() => { + cleanup(); + resolve(); + }, pollMs); + core.once('append', onAppend); + }); } } diff --git a/src/operations/tx/index.js b/src/operations/tx/index.js index 51441aa..f4e758c 100644 --- a/src/operations/tx/index.js +++ b/src/operations/tx/index.js @@ -36,6 +36,11 @@ export class TxOperation { if(false === this.#validator.validate(op)) return; // Stall guard: don't allow a writer to pin apply waiting on an absurd MSB height if (op.value.msbsl > this.#config.maxMsbSignedLength) return; + const localMsbSignedLength = this.#msbClient.getSignedLength(); + if (localMsbSignedLength > 0 && + op.value.msbsl > localMsbSignedLength + this.#config.maxMsbSignedLengthFutureDelta) { + return; + } // Wait for local MSB view to reach the referenced signed length await this.#msbClient.waitForSignedLengthAtLeast(op.value.msbsl); // Fetch MSB apply-op at msbsl by tx key (op.key = tx hash) diff --git a/tests/unit/applyGuards.test.js b/tests/unit/applyGuards.test.js index ac38998..e1f2e3b 100644 --- a/tests/unit/applyGuards.test.js +++ b/tests/unit/applyGuards.test.js @@ -164,6 +164,65 @@ test('apply: tx msbsl stall guard skips waiting', async (t) => { }); }); +test('apply: tx msbsl relative future guard skips waiting', async (t) => { + await withTempDir(async ({ storesDirectory }) => { + const msbBootstrapBuf = b4a.alloc(32).fill(7); + const msb = makeMsbStub({ + msbBootstrapBuf, + signedLength: 100, + async getEntry() { + return null; + }, + }); + + msb.state.base.view.core.once = () => { + throw new Error('apply should not wait for a far-future msbsl'); + }; + + const storeName = 'peer-relative-stall-guard'; + const wallet = await prepareWallet(storesDirectory, storeName); + const config = createConfig(ENV.DEVELOPMENT, { + storesDirectory, + storeName, + maxMsbSignedLength: 1_000_000_000, + maxMsbSignedLengthFutureDelta: 10, + }); + const peer = new Peer({ + config, + msb, + protocol: TestProtocol, + contract: TestContract, + wallet, + }); + + try { + await peer.ready(); + + const txHashHex = makeHex32(1); + const op = { + type: 'tx', + key: txHashHex, + value: { + dispatch: { type: 'ping', value: { msg: 'hi' } }, + msbsl: 111, + ipk: makeHex32(2), + wp: makeHex32(3), + }, + }; + + const timeout = new Promise((_, reject) => + setTimeout(() => reject(new Error('append timed out (possible apply stall)')), 2000) + ); + await Promise.race([peer.base.append(op), timeout]); + + const txl = await peer.bee.get('txl'); + t.is(txl, null, 'tx should not be indexed when msbsl exceeds the local future window'); + } finally { + await closePeer(peer); + } + }); +}); + test('apply: tx MSB payload size guard blocks otherwise-valid tx', async (t) => { await withTempDir(async ({ storesDirectory }) => { const msbBootstrapBuf = b4a.alloc(32).fill(7); From 101e60b5a7a436b42aa7ec120bc0d2e1db72380b Mon Sep 17 00:00:00 2001 From: Benny Date: Mon, 3 Aug 2026 17:16:16 +0200 Subject: [PATCH 02/13] Use Autobase ack for updater acknowledgements --- src/tasks/updater.js | 10 ++++++---- tests/unit/unit.test.js | 1 + tests/unit/updater.test.js | 35 +++++++++++++++++++++++++++++++++++ 3 files changed, 42 insertions(+), 4 deletions(-) create mode 100644 tests/unit/updater.test.js diff --git a/src/tasks/updater.js b/src/tasks/updater.js index ad60bdc..23cb86a 100644 --- a/src/tasks/updater.js +++ b/src/tasks/updater.js @@ -6,9 +6,11 @@ class Updater { #base #scheduler #isInterrupted + #processIntervalMs - constructor({ base }) { + constructor({ base }, config = {}) { this.#base = base + this.#processIntervalMs = Number(config?.updaterIntervalMs ?? PROCESS_INTERVAL_MS) } async start() { @@ -25,18 +27,18 @@ class Updater { async #worker(next) { await this.#update(); - next(PROCESS_INTERVAL_MS); + next(this.#processIntervalMs); } async #update() { if (!this.#shouldRun()) return if (this.#base.view.core.length > this.#base.view.core.signedLength) - await this.#base.append(null) + await this.#base.ack() } #createScheduler() { - return new Scheduler((next) => this.#worker(next), PROCESS_INTERVAL_MS); + return new Scheduler((next) => this.#worker(next), this.#processIntervalMs); } #sleep(ms) { diff --git a/tests/unit/unit.test.js b/tests/unit/unit.test.js index 7e72a38..9063267 100644 --- a/tests/unit/unit.test.js +++ b/tests/unit/unit.test.js @@ -11,4 +11,5 @@ await import('./msbTxValidation.test.js'); await import('./walletNetworkConfig.test.js'); await import('./pearCompat.test.js'); await import('./terminalRuntime.test.js'); +await import('./updater.test.js'); test.resume(); diff --git a/tests/unit/updater.test.js b/tests/unit/updater.test.js new file mode 100644 index 0000000..e173330 --- /dev/null +++ b/tests/unit/updater.test.js @@ -0,0 +1,35 @@ +import test from 'brittle'; +import { Updater } from '../../src/tasks/updater.js'; + +const sleep = (ms) => new Promise(resolve => setTimeout(resolve, ms)); + +test('updater acknowledges unsigned indexer entries without raw null append', async (t) => { + let ackCalls = 0; + let appendCalls = 0; + + const updater = new Updater({ + base: { + isIndexer: true, + view: { + core: { + length: 1, + signedLength: 0 + } + }, + async ack() { + ackCalls++; + }, + async append(value) { + appendCalls++; + throw new Error(`unexpected raw append: ${value}`); + } + } + }, { updaterIntervalMs: 5 }); + + await updater.start(); + await sleep(25); + await updater.stop(); + + t.ok(ackCalls > 0, 'unsigned entries are acknowledged'); + t.is(appendCalls, 0, 'raw null append is not used for acknowledgements'); +}); From bd8a085f7935e204a2964faa860bf212dbdafa49 Mon Sep 17 00:00:00 2001 From: Benny Date: Mon, 3 Aug 2026 17:25:05 +0200 Subject: [PATCH 03/13] Avoid null Autobase ack blocks --- src/index.js | 2 +- src/tasks/updater.js | 10 ++++++++-- tests/unit/updater.test.js | 17 ++++++++--------- 3 files changed, 17 insertions(+), 12 deletions(-) diff --git a/src/index.js b/src/index.js index 0114489..7712600 100644 --- a/src/index.js +++ b/src/index.js @@ -79,7 +79,7 @@ export class Peer extends ReadyResource { async _boot() { this.base = new Autobase(this.store, this.config.bootstrap, { - ackInterval : 1000, + ackInterval : 0, valueEncoding: 'json', open: store => { this.bee = new Hyperbee(store.get('view'), { diff --git a/src/tasks/updater.js b/src/tasks/updater.js index 23cb86a..f3cbf02 100644 --- a/src/tasks/updater.js +++ b/src/tasks/updater.js @@ -1,6 +1,12 @@ import Scheduler from '../utils/scheduler.js'; const PROCESS_INTERVAL_MS = 10_000 +const ACK_OPERATION_TYPE = '_trac_peer_ack_v1' + +const createAckOperation = () => ({ + type: ACK_OPERATION_TYPE, + value: { version: 1 } +}) class Updater { #base @@ -34,7 +40,7 @@ class Updater { if (!this.#shouldRun()) return if (this.#base.view.core.length > this.#base.view.core.signedLength) - await this.#base.ack() + await this.#base.append(createAckOperation()) } #createScheduler() { @@ -56,4 +62,4 @@ class Updater { } } -export { Updater } +export { Updater, ACK_OPERATION_TYPE, createAckOperation } diff --git a/tests/unit/updater.test.js b/tests/unit/updater.test.js index e173330..11b9332 100644 --- a/tests/unit/updater.test.js +++ b/tests/unit/updater.test.js @@ -1,11 +1,10 @@ import test from 'brittle'; -import { Updater } from '../../src/tasks/updater.js'; +import { ACK_OPERATION_TYPE, Updater } from '../../src/tasks/updater.js'; const sleep = (ms) => new Promise(resolve => setTimeout(resolve, ms)); -test('updater acknowledges unsigned indexer entries without raw null append', async (t) => { - let ackCalls = 0; - let appendCalls = 0; +test('updater acknowledges unsigned indexer entries with signed no-op values', async (t) => { + const appended = []; const updater = new Updater({ base: { @@ -17,11 +16,10 @@ test('updater acknowledges unsigned indexer entries without raw null append', as } }, async ack() { - ackCalls++; + throw new Error('autobase null ack path must not be used'); }, async append(value) { - appendCalls++; - throw new Error(`unexpected raw append: ${value}`); + appended.push(value); } } }, { updaterIntervalMs: 5 }); @@ -30,6 +28,7 @@ test('updater acknowledges unsigned indexer entries without raw null append', as await sleep(25); await updater.stop(); - t.ok(ackCalls > 0, 'unsigned entries are acknowledged'); - t.is(appendCalls, 0, 'raw null append is not used for acknowledgements'); + t.ok(appended.length > 0, 'unsigned entries are acknowledged'); + t.ok(appended.every(value => value !== null), 'raw null append is not used for acknowledgements'); + t.is(appended[0].type, ACK_OPERATION_TYPE, 'acknowledgement uses the reserved no-op operation'); }); From de4a589df7a14c6fef727fc2804099ed6a8c8987 Mon Sep 17 00:00:00 2001 From: Benny Date: Mon, 3 Aug 2026 17:36:40 +0200 Subject: [PATCH 04/13] Document no-null Autobase ack patch --- src/index.js | 2 ++ src/tasks/updater.js | 2 ++ 2 files changed, 4 insertions(+) diff --git a/src/index.js b/src/index.js index 7712600..58d23e5 100644 --- a/src/index.js +++ b/src/index.js @@ -79,6 +79,8 @@ export class Peer extends ReadyResource { async _boot() { this.base = new Autobase(this.store, this.config.bootstrap, { + // MAYHEM PATCH: disable Autobase's implicit null ACK; updater appends + // a signed no-op instead so Pear/Bare never encodes a null head. ackInterval : 0, valueEncoding: 'json', open: store => { diff --git a/src/tasks/updater.js b/src/tasks/updater.js index f3cbf02..e71abcf 100644 --- a/src/tasks/updater.js +++ b/src/tasks/updater.js @@ -3,6 +3,8 @@ import Scheduler from '../utils/scheduler.js'; const PROCESS_INTERVAL_MS = 10_000 const ACK_OPERATION_TYPE = '_trac_peer_ack_v1' +// MAYHEM PATCH: ACKs must be ordinary signed operations. append(null) can crash +// under the Pear/Bare hypercore-storage encoder. const createAckOperation = () => ({ type: ACK_OPERATION_TYPE, value: { version: 1 } From 91a383a367e3d5b463b7cb2bec0c935823596a6f Mon Sep 17 00:00:00 2001 From: Benny Date: Mon, 3 Aug 2026 18:13:20 +0200 Subject: [PATCH 05/13] Avoid Autobase null acknowledgements --- src/index.js | 4 +- src/tasks/updater.js | 164 ++++++++++++++++++- tests/unit/updater.test.js | 313 ++++++++++++++++++++++++++++++++++++- 3 files changed, 478 insertions(+), 3 deletions(-) diff --git a/src/index.js b/src/index.js index 58d23e5..24b2a70 100644 --- a/src/index.js +++ b/src/index.js @@ -14,7 +14,7 @@ import { MsbClient } from './msbClient.js'; import { handlerFor } from './operations/index.js'; import TransactionPool from './transaction/transactionPool.js'; import { TransactionObserver } from './tasks/transactionObserver.js'; -import { Updater } from './tasks/updater.js'; +import { Updater, installNonNullAutobaseAck, installSignedAutobaseStore } from './tasks/updater.js'; export { ensureTextCodecs } from './textCodec.js'; export {default as Protocol} from "./artifacts/protocol.js"; export {default as Contract} from "./artifacts/contract.js"; @@ -84,6 +84,7 @@ export class Peer extends ReadyResource { ackInterval : 0, valueEncoding: 'json', open: store => { + installSignedAutobaseStore(store) this.bee = new Hyperbee(store.get('view'), { extension: false, keyEncoding: 'utf-8', @@ -113,6 +114,7 @@ export class Peer extends ReadyResource { await batch.close(); } }) + installNonNullAutobaseAck(this.base) this.base.on('warning', (e) => console.log(e)) } diff --git a/src/tasks/updater.js b/src/tasks/updater.js index e71abcf..5ce1cdf 100644 --- a/src/tasks/updater.js +++ b/src/tasks/updater.js @@ -1,4 +1,5 @@ import Scheduler from '../utils/scheduler.js'; +import b4a from 'b4a'; const PROCESS_INTERVAL_MS = 10_000 const ACK_OPERATION_TYPE = '_trac_peer_ack_v1' @@ -10,6 +11,167 @@ const createAckOperation = () => ({ value: { version: 1 } }) +const canSessionUseKeyPair = (session, keyPair) => { + if (!keyPair?.publicKey) return false + const signers = session.manifest?.signers ?? session.core?.header?.manifest?.signers ?? [] + return signers.some((signer) => b4a.equals(signer.publicKey, keyPair.publicKey)) +} + +const sessionSignerCount = (session) => + (session.manifest?.signers ?? session.core?.header?.manifest?.signers ?? []).length + +const signingAppendOptions = (session, keyPairFor, opts = {}) => { + if (opts?.keyPair || opts?.signature) return opts + const keyPair = keyPairFor() ?? session.keyPair ?? session.core?.header?.keyPair ?? null + return keyPair?.secretKey && canSessionUseKeyPair(session, keyPair) + ? { ...opts, keyPair } + : sessionSignerCount(session) === 0 + ? { ...opts, signature: b4a.alloc(0) } + : opts +} + +const installSignedAutobaseSessionState = (session, keyPairFor) => { + if (!session?.state || typeof session.state.append !== 'function') return session + if (session.state.__tracPeerSignedLocalAppendInstalled === true) return session + + const append = session.state.append.bind(session.state) + Object.defineProperty(session.state, '__tracPeerSignedLocalAppendInstalled', { + value: true, + enumerable: false, + configurable: false + }) + session.state.append = (values, opts = {}) => append(values, signingAppendOptions(session, keyPairFor, opts)) + return session +} + +const installSignedAutobaseSession = (session, keyPairFor) => { + if (!session || typeof session.append !== 'function') return session + if (session.__tracPeerSignedLocalAppendInstalled === true) { + installSignedAutobaseSessionState(session, keyPairFor) + return session + } + + const append = session.append.bind(session) + Object.defineProperty(session, '__tracPeerSignedLocalAppendInstalled', { + value: true, + enumerable: false, + configurable: false + }) + session.append = (blocks, opts = {}) => { + return append(blocks, signingAppendOptions(session, keyPairFor, opts)) + } + installSignedAutobaseSessionState(session, keyPairFor) + if (typeof session.ready === 'function' && + session.__tracPeerSignedReadyInstalled !== true) { + const ready = session.ready.bind(session) + Object.defineProperty(session, '__tracPeerSignedReadyInstalled', { + value: true, + enumerable: false, + configurable: false + }) + session.ready = async (...args) => { + const result = await ready(...args) + installSignedAutobaseSessionState(session, keyPairFor) + return result + } + } + + if (typeof session.session === 'function' && + session.__tracPeerSignedChildSessionInstalled !== true) { + const createChildSession = session.session.bind(session) + Object.defineProperty(session, '__tracPeerSignedChildSessionInstalled', { + value: true, + enumerable: false, + configurable: false + }) + session.session = (...args) => installSignedAutobaseSession(createChildSession(...args), keyPairFor) + } + return session +} + +const installSignedAutobaseView = (view, keyPairFor) => { + if (!view || typeof view.createSession !== 'function') return view + if (view.__tracPeerSignedViewInstalled === true) return view + + const createSession = view.createSession.bind(view) + Object.defineProperty(view, '__tracPeerSignedViewInstalled', { + value: true, + enumerable: false, + configurable: false + }) + view.createSession = (...args) => { + const session = createSession(...args) + installSignedAutobaseSession(view.core, keyPairFor) + installSignedAutobaseSession(view.batch, keyPairFor) + installSignedAutobaseSession(view.atomicBatch, keyPairFor) + return installSignedAutobaseSession(session, keyPairFor) + } + return view +} + +const installSignedAutobaseStore = (store, keyPairFor = null) => { + if (!store) return store + if (store.__tracPeerSignedLocalStoreInstalled === true) return store + + const resolveKeyPair = typeof keyPairFor === 'function' + ? keyPairFor + : () => keyPairFor ?? store.base?.local?.keyPair ?? store.base?.local?.core?.header?.keyPair ?? null + const getLocal = typeof store.getLocal === 'function' + ? store.getLocal.bind(store) + : null + const getViewByName = typeof store.getViewByName === 'function' + ? store.getViewByName.bind(store) + : null + const get = typeof store.get === 'function' + ? store.get.bind(store) + : null + Object.defineProperty(store, '__tracPeerSignedLocalStoreInstalled', { + value: true, + enumerable: false, + configurable: false + }) + if (getLocal) { + store.getLocal = (...args) => installSignedAutobaseSession(getLocal(...args), resolveKeyPair) + } + if (getViewByName) { + store.getViewByName = (...args) => installSignedAutobaseView(getViewByName(...args), resolveKeyPair) + } + if (get) { + store.get = (...args) => installSignedAutobaseSession(get(...args), resolveKeyPair) + } + if (typeof store.atomize === 'function' && + store.__tracPeerSignedAtomizeInstalled !== true) { + const atomize = store.atomize.bind(store) + Object.defineProperty(store, '__tracPeerSignedAtomizeInstalled', { + value: true, + enumerable: false, + configurable: false + }) + store.atomize = (...args) => installSignedAutobaseStore(atomize(...args), resolveKeyPair) + } + return store +} + +const installSignedLocalAutobaseAppend = (base) => { + if (!base?._viewStore) return base + return installSignedAutobaseStore(base._viewStore, () => base.local?.keyPair ?? base.local?.core?.header?.keyPair ?? null) +} + +const installNonNullAutobaseAck = (base) => { + if (!base || typeof base.append !== 'function') return base + if (base.__tracPeerNonNullAckInstalled === true) return base + + const append = base.append.bind(base) + Object.defineProperty(base, '__tracPeerNonNullAckInstalled', { + value: true, + enumerable: false, + configurable: false + }) + base.append = (value, ...args) => append(value === null ? createAckOperation() : value, ...args) + installSignedLocalAutobaseAppend(base) + return base +} + class Updater { #base #scheduler @@ -64,4 +226,4 @@ class Updater { } } -export { Updater, ACK_OPERATION_TYPE, createAckOperation } +export { Updater, ACK_OPERATION_TYPE, createAckOperation, installNonNullAutobaseAck, installSignedAutobaseStore } diff --git a/tests/unit/updater.test.js b/tests/unit/updater.test.js index 11b9332..2e0334b 100644 --- a/tests/unit/updater.test.js +++ b/tests/unit/updater.test.js @@ -1,5 +1,5 @@ import test from 'brittle'; -import { ACK_OPERATION_TYPE, Updater } from '../../src/tasks/updater.js'; +import { ACK_OPERATION_TYPE, Updater, installNonNullAutobaseAck, installSignedAutobaseStore } from '../../src/tasks/updater.js'; const sleep = (ms) => new Promise(resolve => setTimeout(resolve, ms)); @@ -32,3 +32,314 @@ test('updater acknowledges unsigned indexer entries with signed no-op values', a t.ok(appended.every(value => value !== null), 'raw null append is not used for acknowledgements'); t.is(appended[0].type, ACK_OPERATION_TYPE, 'acknowledgement uses the reserved no-op operation'); }); + +test('installNonNullAutobaseAck maps Autobase null acknowledgements to signed no-op values', async (t) => { + const appended = []; + const base = { + async append(value) { + appended.push(value); + return value; + } + }; + + installNonNullAutobaseAck(base); + installNonNullAutobaseAck(base); + + const ack = await base.append(null); + const payload = { type: 'msg', value: 'hello' }; + const passthrough = await base.append(payload); + + t.is(appended.length, 2, 'append is wrapped once'); + t.ok(ack !== null, 'null acknowledgement is replaced'); + t.is(ack.type, ACK_OPERATION_TYPE, 'replacement uses reserved ack operation'); + t.is(appended[0].type, ACK_OPERATION_TYPE, 'stored append value is non-null ack operation'); + t.is(passthrough, payload, 'non-null appends are unchanged'); + t.is(appended[1], payload, 'non-null stored value is unchanged'); +}); + +test('installNonNullAutobaseAck signs Autobase local named-session appends', async (t) => { + const keyPair = { publicKey: Buffer.alloc(32, 1), secretKey: Buffer.alloc(64, 2) }; + const appendCalls = []; + const local = { + manifest: { signers: [{ publicKey: keyPair.publicKey }] }, + async append(blocks, opts = {}) { + appendCalls.push({ blocks, opts }); + } + }; + const localStore = { + getLocal() { + return local; + } + }; + const rootStore = { + getLocal() { + return local; + }, + atomize() { + return localStore; + } + }; + const base = { + local: { keyPair }, + _viewStore: rootStore, + async append(value) { + return value; + } + }; + + installNonNullAutobaseAck(base); + + await rootStore.getLocal().append(['root']); + await rootStore.atomize().getLocal().append(['local']); + const explicit = { signature: Buffer.alloc(64, 3) }; + await rootStore.getLocal().append(['explicit'], explicit); + + t.is(appendCalls.length, 3, 'local append wrappers are installed'); + t.is(appendCalls[0].opts.keyPair, keyPair, 'root local append receives the writer keypair'); + t.is(appendCalls[1].opts.keyPair, keyPair, 'atomized local append receives the writer keypair'); + t.is(appendCalls[2].opts, explicit, 'explicit signing options are preserved'); +}); + +test('installNonNullAutobaseAck signs Autobase view batch session appends', async (t) => { + const keyPair = { publicKey: Buffer.alloc(32, 4), secretKey: Buffer.alloc(64, 5) }; + const appendCalls = []; + const makeSession = () => ({ + manifest: { signers: [{ publicKey: keyPair.publicKey }] }, + async append(blocks, opts = {}) { + appendCalls.push({ blocks, opts }); + }, + session() { + return makeSession(); + } + }); + const view = { + core: makeSession(), + batch: null, + atomicBatch: null, + createSession() { + this.batch = makeSession(); + return this.batch.session(); + } + }; + const rootStore = { + getLocal() { + return makeSession(); + }, + getViewByName() { + return view; + }, + get(name) { + return this.getViewByName(name).createSession(); + }, + atomize() { + return this; + } + }; + const base = { + local: { keyPair }, + _viewStore: rootStore, + async append(value) { + return value; + } + }; + + installNonNullAutobaseAck(base); + + const viewSession = rootStore.get('view'); + await view.batch.append(['batch']); + await viewSession.append(['child']); + + t.is(appendCalls.length, 2, 'view batch and child session append were captured'); + t.is(appendCalls[0].opts.keyPair, keyPair, 'view batch append receives the writer keypair'); + t.is(appendCalls[1].opts.keyPair, keyPair, 'view child append receives the writer keypair'); +}); + +test('installSignedAutobaseStore signs view sessions created during Autobase open', async (t) => { + const keyPair = { publicKey: Buffer.alloc(32, 6), secretKey: Buffer.alloc(64, 7) }; + const appendCalls = []; + const base = { local: null }; + const makeSession = () => ({ + manifest: { signers: [{ publicKey: keyPair.publicKey }] }, + async append(blocks, opts = {}) { + appendCalls.push({ blocks, opts }); + }, + session() { + return makeSession(); + } + }); + const view = { + core: makeSession(), + batch: null, + atomicBatch: null, + createSession() { + this.batch = makeSession(); + return this.batch.session(); + } + }; + const store = { + base, + getLocal() { + return makeSession(); + }, + getViewByName() { + return view; + }, + get(name) { + return this.getViewByName(name).createSession(); + }, + atomize() { + return this; + } + }; + + installSignedAutobaseStore(store); + const viewSession = store.get('view'); + base.local = { keyPair }; + + await view.batch.append(['batch']); + await viewSession.append(['child']); + + t.is(appendCalls.length, 2, 'open-created view sessions are wrapped'); + t.is(appendCalls[0].opts.keyPair, keyPair, 'open-created batch append receives the writer keypair'); + t.is(appendCalls[1].opts.keyPair, keyPair, 'open-created child append receives the writer keypair'); +}); + +test('installSignedAutobaseStore uses empty signatures for unsigned genesis views', async (t) => { + const appendCalls = []; + const makeSession = () => ({ + manifest: { signers: [] }, + core: { header: { manifest: { signers: [] } } }, + async append(blocks, opts = {}) { + appendCalls.push({ blocks, opts }); + }, + session() { + return makeSession(); + } + }); + const view = { + core: makeSession(), + batch: null, + atomicBatch: null, + createSession() { + this.batch = makeSession(); + return this.batch.session(); + } + }; + const store = { + getLocal() { + return makeSession(); + }, + getViewByName() { + return view; + }, + get(name) { + return this.getViewByName(name).createSession(); + } + }; + + installSignedAutobaseStore(store); + const viewSession = store.get('view'); + await view.batch.append(['genesis-batch']); + await viewSession.append(['genesis-child']); + + t.is(appendCalls.length, 2, 'unsigned genesis appends are captured'); + t.is(appendCalls[0].opts.signature.length, 0, 'genesis batch append receives an empty signature buffer'); + t.is(appendCalls[1].opts.signature.length, 0, 'genesis child append receives an empty signature buffer'); +}); + +test('installSignedAutobaseStore covers direct session-state genesis appends', async (t) => { + const appendCalls = []; + const makeSession = () => ({ + manifest: { signers: [] }, + core: { header: { manifest: { signers: [] } } }, + state: { + async append(values, opts = {}) { + appendCalls.push({ values, opts }); + } + }, + async append() { + throw new Error('test must exercise direct state append'); + }, + session() { + return makeSession(); + } + }); + const view = { + core: makeSession(), + batch: null, + atomicBatch: null, + createSession() { + this.batch = makeSession(); + return this.batch.session(); + } + }; + const store = { + getLocal() { + return makeSession(); + }, + getViewByName() { + return view; + }, + get(name) { + return this.getViewByName(name).createSession(); + } + }; + + installSignedAutobaseStore(store); + const viewSession = store.get('view'); + + await viewSession.state.append(['encoded-genesis']); + + t.is(appendCalls.length, 1, 'direct state append is wrapped'); + t.is(appendCalls[0].opts.signature.length, 0, 'direct state append receives an empty signature buffer'); +}); + +test('installSignedAutobaseStore installs state append wrapper after session ready', async (t) => { + const appendCalls = []; + const makeSession = () => ({ + manifest: { signers: [] }, + core: { header: { manifest: { signers: [] } } }, + state: null, + async ready() { + this.state = { + async append(values, opts = {}) { + appendCalls.push({ values, opts }); + } + }; + }, + async append() { + throw new Error('test must exercise ready-installed state append'); + }, + session() { + return makeSession(); + } + }); + const view = { + core: makeSession(), + batch: null, + atomicBatch: null, + createSession() { + this.batch = makeSession(); + return this.batch.session(); + } + }; + const store = { + getLocal() { + return makeSession(); + }, + getViewByName() { + return view; + }, + get(name) { + return this.getViewByName(name).createSession(); + } + }; + + installSignedAutobaseStore(store); + const viewSession = store.get('view'); + await viewSession.ready(); + await viewSession.state.append(['encoded-after-ready']); + + t.is(appendCalls.length, 1, 'state append created during ready is wrapped'); + t.is(appendCalls[0].opts.signature.length, 0, 'ready-created state append receives an empty signature buffer'); +}); From 8c8af52a4af61e6f8b06f4a953ab5460747bda17 Mon Sep 17 00:00:00 2001 From: Benny Date: Mon, 3 Aug 2026 18:52:45 +0200 Subject: [PATCH 06/13] Normalize migrated Autobase view heads --- src/tasks/updater.js | 41 ++++++++++++++++++++++++++++++++++++++ tests/unit/updater.test.js | 39 ++++++++++++++++++++++++++++++++++++ 2 files changed, 80 insertions(+) diff --git a/src/tasks/updater.js b/src/tasks/updater.js index 5ce1cdf..641c57a 100644 --- a/src/tasks/updater.js +++ b/src/tasks/updater.js @@ -30,6 +30,39 @@ const signingAppendOptions = (session, keyPairFor, opts = {}) => { : opts } +const normalizeAutobaseHead = (head) => + head && (head.signature === null || head.signature === undefined) + ? { ...head, signature: b4a.alloc(0) } + : head + +const installSignedAutobaseCoreStorage = (core) => { + if (!core?.storage || typeof core.storage.write !== 'function') return core + if (core.storage.__tracPeerSignedCoreTxInstalled === true) return core + + const write = core.storage.write.bind(core.storage) + Object.defineProperty(core.storage, '__tracPeerSignedCoreTxInstalled', { + value: true, + enumerable: false, + configurable: false + }) + core.storage.write = (...args) => { + const tx = write(...args) + if (!tx || typeof tx.setHead !== 'function' || + tx.__tracPeerSignedCoreTxHeadInstalled === true) { + return tx + } + const setHead = tx.setHead.bind(tx) + Object.defineProperty(tx, '__tracPeerSignedCoreTxHeadInstalled', { + value: true, + enumerable: false, + configurable: false + }) + tx.setHead = (head, ...headArgs) => setHead(normalizeAutobaseHead(head), ...headArgs) + return tx + } + return core +} + const installSignedAutobaseSessionState = (session, keyPairFor) => { if (!session?.state || typeof session.state.append !== 'function') return session if (session.state.__tracPeerSignedLocalAppendInstalled === true) return session @@ -46,6 +79,7 @@ const installSignedAutobaseSessionState = (session, keyPairFor) => { const installSignedAutobaseSession = (session, keyPairFor) => { if (!session || typeof session.append !== 'function') return session + installSignedAutobaseCoreStorage(session.core) if (session.__tracPeerSignedLocalAppendInstalled === true) { installSignedAutobaseSessionState(session, keyPairFor) return session @@ -71,6 +105,7 @@ const installSignedAutobaseSession = (session, keyPairFor) => { }) session.ready = async (...args) => { const result = await ready(...args) + installSignedAutobaseCoreStorage(session.core) installSignedAutobaseSessionState(session, keyPairFor) return result } @@ -122,6 +157,9 @@ const installSignedAutobaseStore = (store, keyPairFor = null) => { const getViewByName = typeof store.getViewByName === 'function' ? store.getViewByName.bind(store) : null + const getViewCore = typeof store.getViewCore === 'function' + ? store.getViewCore.bind(store) + : null const get = typeof store.get === 'function' ? store.get.bind(store) : null @@ -136,6 +174,9 @@ const installSignedAutobaseStore = (store, keyPairFor = null) => { if (getViewByName) { store.getViewByName = (...args) => installSignedAutobaseView(getViewByName(...args), resolveKeyPair) } + if (getViewCore) { + store.getViewCore = (...args) => installSignedAutobaseSession(getViewCore(...args), resolveKeyPair) + } if (get) { store.get = (...args) => installSignedAutobaseSession(get(...args), resolveKeyPair) } diff --git a/tests/unit/updater.test.js b/tests/unit/updater.test.js index 2e0334b..a93e240 100644 --- a/tests/unit/updater.test.js +++ b/tests/unit/updater.test.js @@ -343,3 +343,42 @@ test('installSignedAutobaseStore installs state append wrapper after session rea t.is(appendCalls.length, 1, 'state append created during ready is wrapped'); t.is(appendCalls[0].opts.signature.length, 0, 'ready-created state append receives an empty signature buffer'); }); + +test('installSignedAutobaseStore covers Autobase migrated view cores', async (t) => { + const heads = []; + const viewCoreSession = { + manifest: { signers: [] }, + core: { + header: { manifest: { signers: [] } }, + storage: { + write() { + return { + setHead(head) { + heads.push(head); + } + }; + } + } + }, + async ready() {}, + async append() {} + }; + const store = { + getViewCore() { + return viewCoreSession; + } + }; + + installSignedAutobaseStore(store); + const migrated = store.getViewCore(); + await migrated.ready(); + migrated.core.storage.write().setHead({ + fork: 0, + length: 1, + rootHash: Buffer.alloc(32, 8), + signature: null + }); + + t.is(heads.length, 1, 'migrated view core storage transaction is wrapped'); + t.is(heads[0].signature.length, 0, 'null prologue head signature becomes an empty signature buffer'); +}); From 05d58408783cf4aef733e095bf6a99f1d94405c1 Mon Sep 17 00:00:00 2001 From: Benny Date: Mon, 3 Aug 2026 18:57:54 +0200 Subject: [PATCH 07/13] Normalize Autobase session heads --- src/tasks/updater.js | 7 +++++++ tests/unit/updater.test.js | 39 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 46 insertions(+) diff --git a/src/tasks/updater.js b/src/tasks/updater.js index 641c57a..0195c5a 100644 --- a/src/tasks/updater.js +++ b/src/tasks/updater.js @@ -40,11 +40,18 @@ const installSignedAutobaseCoreStorage = (core) => { if (core.storage.__tracPeerSignedCoreTxInstalled === true) return core const write = core.storage.write.bind(core.storage) + const createSession = typeof core.storage.createSession === 'function' + ? core.storage.createSession.bind(core.storage) + : null Object.defineProperty(core.storage, '__tracPeerSignedCoreTxInstalled', { value: true, enumerable: false, configurable: false }) + if (createSession) { + core.storage.createSession = (name, head, ...args) => + createSession(name, normalizeAutobaseHead(head), ...args) + } core.storage.write = (...args) => { const tx = write(...args) if (!tx || typeof tx.setHead !== 'function' || diff --git a/tests/unit/updater.test.js b/tests/unit/updater.test.js index a93e240..6927990 100644 --- a/tests/unit/updater.test.js +++ b/tests/unit/updater.test.js @@ -382,3 +382,42 @@ test('installSignedAutobaseStore covers Autobase migrated view cores', async (t) t.is(heads.length, 1, 'migrated view core storage transaction is wrapped'); t.is(heads[0].signature.length, 0, 'null prologue head signature becomes an empty signature buffer'); }); + +test('installSignedAutobaseStore covers named session head creation', async (t) => { + const heads = []; + const viewCoreSession = { + manifest: { signers: [] }, + core: { + header: { manifest: { signers: [] } }, + storage: { + write() { + return { setHead() {} }; + }, + async createSession(name, head) { + heads.push({ name, head }); + return {}; + } + } + }, + async ready() {}, + async append() {} + }; + const store = { + getViewCore() { + return viewCoreSession; + } + }; + + installSignedAutobaseStore(store); + const migrated = store.getViewCore(); + await migrated.ready(); + await migrated.core.storage.createSession('batch', { + fork: 0, + length: 1, + rootHash: Buffer.alloc(32, 9), + signature: null + }); + + t.is(heads.length, 1, 'named session creation is wrapped'); + t.is(heads[0].head.signature.length, 0, 'null session head signature becomes an empty signature buffer'); +}); From b2d22b4c933b5a66e55f4c0fe73026bad136b55a Mon Sep 17 00:00:00 2001 From: Benny Date: Mon, 3 Aug 2026 19:03:58 +0200 Subject: [PATCH 08/13] Normalize Autobase atomic session heads --- src/tasks/updater.js | 7 +++++++ tests/unit/updater.test.js | 39 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 46 insertions(+) diff --git a/src/tasks/updater.js b/src/tasks/updater.js index 0195c5a..094925a 100644 --- a/src/tasks/updater.js +++ b/src/tasks/updater.js @@ -43,6 +43,9 @@ const installSignedAutobaseCoreStorage = (core) => { const createSession = typeof core.storage.createSession === 'function' ? core.storage.createSession.bind(core.storage) : null + const createAtomicSession = typeof core.storage.createAtomicSession === 'function' + ? core.storage.createAtomicSession.bind(core.storage) + : null Object.defineProperty(core.storage, '__tracPeerSignedCoreTxInstalled', { value: true, enumerable: false, @@ -52,6 +55,10 @@ const installSignedAutobaseCoreStorage = (core) => { core.storage.createSession = (name, head, ...args) => createSession(name, normalizeAutobaseHead(head), ...args) } + if (createAtomicSession) { + core.storage.createAtomicSession = (atom, head, ...args) => + createAtomicSession(atom, normalizeAutobaseHead(head), ...args) + } core.storage.write = (...args) => { const tx = write(...args) if (!tx || typeof tx.setHead !== 'function' || diff --git a/tests/unit/updater.test.js b/tests/unit/updater.test.js index 6927990..bc9627e 100644 --- a/tests/unit/updater.test.js +++ b/tests/unit/updater.test.js @@ -421,3 +421,42 @@ test('installSignedAutobaseStore covers named session head creation', async (t) t.is(heads.length, 1, 'named session creation is wrapped'); t.is(heads[0].head.signature.length, 0, 'null session head signature becomes an empty signature buffer'); }); + +test('installSignedAutobaseStore covers atomic session head creation', async (t) => { + const heads = []; + const viewCoreSession = { + manifest: { signers: [] }, + core: { + header: { manifest: { signers: [] } }, + storage: { + write() { + return { setHead() {} }; + }, + async createAtomicSession(atom, head) { + heads.push({ atom, head }); + return {}; + } + } + }, + async ready() {}, + async append() {} + }; + const store = { + getViewCore() { + return viewCoreSession; + } + }; + + installSignedAutobaseStore(store); + const migrated = store.getViewCore(); + await migrated.ready(); + await migrated.core.storage.createAtomicSession({ view: {} }, { + fork: 0, + length: 1, + rootHash: Buffer.alloc(32, 10), + signature: null + }); + + t.is(heads.length, 1, 'atomic session creation is wrapped'); + t.is(heads[0].head.signature.length, 0, 'null atomic session head signature becomes an empty signature buffer'); +}); From 5603e06c25ae5bdd0b7ebcc9517add44a769178d Mon Sep 17 00:00:00 2001 From: Benny Date: Mon, 3 Aug 2026 19:13:09 +0200 Subject: [PATCH 09/13] Install Autobase storage guard before boot --- src/index.js | 1 + 1 file changed, 1 insertion(+) diff --git a/src/index.js b/src/index.js index 24b2a70..3b96fda 100644 --- a/src/index.js +++ b/src/index.js @@ -78,6 +78,7 @@ export class Peer extends ReadyResource { } async _boot() { + installSignedAutobaseStore(this.store, () => this.base?.local?.keyPair ?? null) this.base = new Autobase(this.store, this.config.bootstrap, { // MAYHEM PATCH: disable Autobase's implicit null ACK; updater appends // a signed no-op instead so Pear/Bare never encodes a null head. From f11d6a12ec27342303a6a2eee72530c150d0a8e2 Mon Sep 17 00:00:00 2001 From: Benny Date: Mon, 3 Aug 2026 19:17:52 +0200 Subject: [PATCH 10/13] Normalize Autobase state storage heads --- src/tasks/updater.js | 31 +++++++++++++++---------- tests/unit/updater.test.js | 46 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 65 insertions(+), 12 deletions(-) diff --git a/src/tasks/updater.js b/src/tasks/updater.js index 094925a..b0f0c3e 100644 --- a/src/tasks/updater.js +++ b/src/tasks/updater.js @@ -35,31 +35,31 @@ const normalizeAutobaseHead = (head) => ? { ...head, signature: b4a.alloc(0) } : head -const installSignedAutobaseCoreStorage = (core) => { - if (!core?.storage || typeof core.storage.write !== 'function') return core - if (core.storage.__tracPeerSignedCoreTxInstalled === true) return core +const installSignedAutobaseStorage = (storage) => { + if (!storage || typeof storage.write !== 'function') return storage + if (storage.__tracPeerSignedCoreTxInstalled === true) return storage - const write = core.storage.write.bind(core.storage) - const createSession = typeof core.storage.createSession === 'function' - ? core.storage.createSession.bind(core.storage) + const write = storage.write.bind(storage) + const createSession = typeof storage.createSession === 'function' + ? storage.createSession.bind(storage) : null - const createAtomicSession = typeof core.storage.createAtomicSession === 'function' - ? core.storage.createAtomicSession.bind(core.storage) + const createAtomicSession = typeof storage.createAtomicSession === 'function' + ? storage.createAtomicSession.bind(storage) : null - Object.defineProperty(core.storage, '__tracPeerSignedCoreTxInstalled', { + Object.defineProperty(storage, '__tracPeerSignedCoreTxInstalled', { value: true, enumerable: false, configurable: false }) if (createSession) { - core.storage.createSession = (name, head, ...args) => + storage.createSession = (name, head, ...args) => createSession(name, normalizeAutobaseHead(head), ...args) } if (createAtomicSession) { - core.storage.createAtomicSession = (atom, head, ...args) => + storage.createAtomicSession = (atom, head, ...args) => createAtomicSession(atom, normalizeAutobaseHead(head), ...args) } - core.storage.write = (...args) => { + storage.write = (...args) => { const tx = write(...args) if (!tx || typeof tx.setHead !== 'function' || tx.__tracPeerSignedCoreTxHeadInstalled === true) { @@ -74,6 +74,13 @@ const installSignedAutobaseCoreStorage = (core) => { tx.setHead = (head, ...headArgs) => setHead(normalizeAutobaseHead(head), ...headArgs) return tx } + return storage +} + +const installSignedAutobaseCoreStorage = (core) => { + if (!core) return core + installSignedAutobaseStorage(core.storage) + installSignedAutobaseStorage(core.state?.storage) return core } diff --git a/tests/unit/updater.test.js b/tests/unit/updater.test.js index bc9627e..c640064 100644 --- a/tests/unit/updater.test.js +++ b/tests/unit/updater.test.js @@ -460,3 +460,49 @@ test('installSignedAutobaseStore covers atomic session head creation', async (t) t.is(heads.length, 1, 'atomic session creation is wrapped'); t.is(heads[0].head.signature.length, 0, 'null atomic session head signature becomes an empty signature buffer'); }); + +test('installSignedAutobaseStore covers core state atomic session heads', async (t) => { + const heads = []; + const viewCoreSession = { + manifest: { signers: [] }, + core: { + header: { manifest: { signers: [] } }, + storage: { + write() { + return { setHead() {} }; + } + }, + state: { + storage: { + write() { + return { setHead() {} }; + }, + async createAtomicSession(atom, head) { + heads.push({ atom, head }); + return {}; + } + } + } + }, + async ready() {}, + async append() {} + }; + const store = { + getViewCore() { + return viewCoreSession; + } + }; + + installSignedAutobaseStore(store); + const migrated = store.getViewCore(); + await migrated.ready(); + await migrated.core.state.storage.createAtomicSession({ view: {} }, { + fork: 0, + length: 1, + rootHash: Buffer.alloc(32, 11), + signature: null + }); + + t.is(heads.length, 1, 'core state atomic session creation is wrapped'); + t.is(heads[0].head.signature.length, 0, 'null state atomic session head signature becomes an empty signature buffer'); +}); From 0856b36c43e38f5ba88f08d37832d2c28d182ce9 Mon Sep 17 00:00:00 2001 From: Benny Date: Mon, 3 Aug 2026 19:25:05 +0200 Subject: [PATCH 11/13] Normalize Corestore-created storage heads --- src/index.js | 3 ++- src/tasks/updater.js | 25 +++++++++++++++++++++++- tests/unit/updater.test.js | 40 +++++++++++++++++++++++++++++++++++++- 3 files changed, 65 insertions(+), 3 deletions(-) diff --git a/src/index.js b/src/index.js index 3b96fda..9da5f5b 100644 --- a/src/index.js +++ b/src/index.js @@ -14,7 +14,7 @@ import { MsbClient } from './msbClient.js'; import { handlerFor } from './operations/index.js'; import TransactionPool from './transaction/transactionPool.js'; import { TransactionObserver } from './tasks/transactionObserver.js'; -import { Updater, installNonNullAutobaseAck, installSignedAutobaseStore } from './tasks/updater.js'; +import { Updater, installNonNullAutobaseAck, installSignedAutobaseStore, installSignedCorestoreStorageFactory } from './tasks/updater.js'; export { ensureTextCodecs } from './textCodec.js'; export {default as Protocol} from "./artifacts/protocol.js"; export {default as Contract} from "./artifacts/contract.js"; @@ -78,6 +78,7 @@ export class Peer extends ReadyResource { } async _boot() { + installSignedCorestoreStorageFactory(this.store) installSignedAutobaseStore(this.store, () => this.base?.local?.keyPair ?? null) this.base = new Autobase(this.store, this.config.bootstrap, { // MAYHEM PATCH: disable Autobase's implicit null ACK; updater appends diff --git a/src/tasks/updater.js b/src/tasks/updater.js index b0f0c3e..80cc07c 100644 --- a/src/tasks/updater.js +++ b/src/tasks/updater.js @@ -84,6 +84,29 @@ const installSignedAutobaseCoreStorage = (core) => { return core } +const installSignedCorestoreStorageFactory = (store) => { + const storage = store?.storage ?? store + if (!storage || storage.__tracPeerSignedStorageFactoryInstalled === true) return store + + Object.defineProperty(storage, '__tracPeerSignedStorageFactoryInstalled', { + value: true, + enumerable: false, + configurable: false + }) + + for (const method of ['create', 'resume', '_create', '_resumeFromPointers']) { + if (typeof storage[method] !== 'function') continue + const original = storage[method].bind(storage) + storage[method] = (...args) => { + const result = original(...args) + return result && typeof result.then === 'function' + ? result.then(installSignedAutobaseStorage) + : installSignedAutobaseStorage(result) + } + } + return store +} + const installSignedAutobaseSessionState = (session, keyPairFor) => { if (!session?.state || typeof session.state.append !== 'function') return session if (session.state.__tracPeerSignedLocalAppendInstalled === true) return session @@ -288,4 +311,4 @@ class Updater { } } -export { Updater, ACK_OPERATION_TYPE, createAckOperation, installNonNullAutobaseAck, installSignedAutobaseStore } +export { Updater, ACK_OPERATION_TYPE, createAckOperation, installNonNullAutobaseAck, installSignedAutobaseStore, installSignedCorestoreStorageFactory } diff --git a/tests/unit/updater.test.js b/tests/unit/updater.test.js index c640064..7d8ad8b 100644 --- a/tests/unit/updater.test.js +++ b/tests/unit/updater.test.js @@ -1,5 +1,11 @@ import test from 'brittle'; -import { ACK_OPERATION_TYPE, Updater, installNonNullAutobaseAck, installSignedAutobaseStore } from '../../src/tasks/updater.js'; +import { + ACK_OPERATION_TYPE, + Updater, + installNonNullAutobaseAck, + installSignedAutobaseStore, + installSignedCorestoreStorageFactory +} from '../../src/tasks/updater.js'; const sleep = (ms) => new Promise(resolve => setTimeout(resolve, ms)); @@ -506,3 +512,35 @@ test('installSignedAutobaseStore covers core state atomic session heads', async t.is(heads.length, 1, 'core state atomic session creation is wrapped'); t.is(heads[0].head.signature.length, 0, 'null state atomic session head signature becomes an empty signature buffer'); }); + +test('installSignedCorestoreStorageFactory wraps created HypercoreStorage instances', async (t) => { + const heads = []; + const hypercoreStorage = { + write() { + return { setHead() {} }; + }, + async createAtomicSession(atom, head) { + heads.push({ atom, head }); + return {}; + } + }; + const corestore = { + storage: { + async create() { + return hypercoreStorage; + } + } + }; + + installSignedCorestoreStorageFactory(corestore); + const storage = await corestore.storage.create(); + await storage.createAtomicSession({ view: {} }, { + fork: 0, + length: 1, + rootHash: Buffer.alloc(32, 12), + signature: null + }); + + t.is(heads.length, 1, 'created HypercoreStorage instance is wrapped'); + t.is(heads[0].head.signature.length, 0, 'factory-normalized atomic head has an empty signature buffer'); +}); From 25d2700f062476f46fa2197c12e2be794def6cb3 Mon Sep 17 00:00:00 2001 From: Benny Date: Mon, 3 Aug 2026 19:33:41 +0200 Subject: [PATCH 12/13] Normalize atomized storage heads --- src/tasks/updater.js | 20 +++++++++++++++++-- tests/unit/updater.test.js | 41 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 59 insertions(+), 2 deletions(-) diff --git a/src/tasks/updater.js b/src/tasks/updater.js index 80cc07c..40a337d 100644 --- a/src/tasks/updater.js +++ b/src/tasks/updater.js @@ -46,6 +46,16 @@ const installSignedAutobaseStorage = (storage) => { const createAtomicSession = typeof storage.createAtomicSession === 'function' ? storage.createAtomicSession.bind(storage) : null + const resumeSession = typeof storage.resumeSession === 'function' + ? storage.resumeSession.bind(storage) + : null + const atomize = typeof storage.atomize === 'function' + ? storage.atomize.bind(storage) + : null + const wrapStorageResult = (result) => + result && typeof result.then === 'function' + ? result.then(installSignedAutobaseStorage) + : installSignedAutobaseStorage(result) Object.defineProperty(storage, '__tracPeerSignedCoreTxInstalled', { value: true, enumerable: false, @@ -53,11 +63,17 @@ const installSignedAutobaseStorage = (storage) => { }) if (createSession) { storage.createSession = (name, head, ...args) => - createSession(name, normalizeAutobaseHead(head), ...args) + wrapStorageResult(createSession(name, normalizeAutobaseHead(head), ...args)) } if (createAtomicSession) { storage.createAtomicSession = (atom, head, ...args) => - createAtomicSession(atom, normalizeAutobaseHead(head), ...args) + wrapStorageResult(createAtomicSession(atom, normalizeAutobaseHead(head), ...args)) + } + if (resumeSession) { + storage.resumeSession = (...args) => wrapStorageResult(resumeSession(...args)) + } + if (atomize) { + storage.atomize = (...args) => wrapStorageResult(atomize(...args)) } storage.write = (...args) => { const tx = write(...args) diff --git a/tests/unit/updater.test.js b/tests/unit/updater.test.js index 7d8ad8b..cccf645 100644 --- a/tests/unit/updater.test.js +++ b/tests/unit/updater.test.js @@ -544,3 +544,44 @@ test('installSignedCorestoreStorageFactory wraps created HypercoreStorage instan t.is(heads.length, 1, 'created HypercoreStorage instance is wrapped'); t.is(heads[0].head.signature.length, 0, 'factory-normalized atomic head has an empty signature buffer'); }); + +test('installSignedCorestoreStorageFactory wraps atomized HypercoreStorage instances', async (t) => { + const heads = []; + const atomizedStorage = { + write() { + return { setHead() {} }; + }, + async createAtomicSession(atom, head) { + heads.push({ atom, head }); + return {}; + } + }; + const hypercoreStorage = { + write() { + return { setHead() {} }; + }, + atomize() { + return atomizedStorage; + } + }; + const corestore = { + storage: { + async create() { + return hypercoreStorage; + } + } + }; + + installSignedCorestoreStorageFactory(corestore); + const storage = await corestore.storage.create(); + const atomized = storage.atomize({ view: {} }); + await atomized.createAtomicSession({ view: {} }, { + fork: 0, + length: 1, + rootHash: Buffer.alloc(32, 13), + signature: null + }); + + t.is(heads.length, 1, 'atomized HypercoreStorage instance is wrapped'); + t.is(heads[0].head.signature.length, 0, 'atomized storage normalizes atomic head signatures'); +}); From d687061cb2016e5c733abb65fb7b3a83c0b32634 Mon Sep 17 00:00:00 2001 From: Benny Date: Thu, 6 Aug 2026 15:21:53 +0200 Subject: [PATCH 13/13] Update protobufjs lock --- package-lock.json | 50 ++++++++++++++++++++--------------------------- 1 file changed, 21 insertions(+), 29 deletions(-) diff --git a/package-lock.json b/package-lock.json index bcd0fad..9db501c 100644 --- a/package-lock.json +++ b/package-lock.json @@ -304,25 +304,24 @@ "license": "BSD-3-Clause" }, "node_modules/@protobufjs/codegen": { - "version": "2.0.4", - "resolved": "https://registry.npmjs.org/@protobufjs/codegen/-/codegen-2.0.4.tgz", - "integrity": "sha512-YyFaikqM5sH0ziFZCN3xDC7zeGaB/d0IUb9CATugHWbd1FRFwWwt4ld4OYMPWu5a3Xe01mGAULCdqhMlPl29Jg==", + "version": "2.0.5", + "resolved": "https://registry.npmjs.org/@protobufjs/codegen/-/codegen-2.0.5.tgz", + "integrity": "sha512-zgXFLzW3Ap33e6d0Wlj4MGIm6Ce8O89n/apUaGNB/jx+hw+ruWEp7EwGUshdLKVRCxZW12fp9r40E1mQrf/34g==", "license": "BSD-3-Clause" }, "node_modules/@protobufjs/eventemitter": { - "version": "1.1.0", - "resolved": "https://registry.npmjs.org/@protobufjs/eventemitter/-/eventemitter-1.1.0.tgz", - "integrity": "sha512-j9ednRT81vYJ9OfVuXG6ERSTdEL1xVsNgqpkxMsbIabzSo3goCjDIveeGv5d03om39ML71RdmrGNjG5SReBP/Q==", + "version": "1.1.1", + "resolved": "https://registry.npmjs.org/@protobufjs/eventemitter/-/eventemitter-1.1.1.tgz", + "integrity": "sha512-vW1GmwMZNnL+gMRaovlh9yZX74kc+TTU3FObkkurpMaRtBfLP3ldjS9KQWlwZgraRE0+dheEEoAxdzcJQ8eXZg==", "license": "BSD-3-Clause" }, "node_modules/@protobufjs/fetch": { - "version": "1.1.0", - "resolved": "https://registry.npmjs.org/@protobufjs/fetch/-/fetch-1.1.0.tgz", - "integrity": "sha512-lljVXpqXebpsijW71PZaCYeIcE5on1w5DlQy5WH6GLbFryLUrBD4932W/E2BSpfRJWseIL4v/KPgBFxDOIdKpQ==", + "version": "1.1.1", + "resolved": "https://registry.npmjs.org/@protobufjs/fetch/-/fetch-1.1.1.tgz", + "integrity": "sha512-GpptLrs57adMSuHi3VNj0mAF8dwh36LMaYF6XyJ6JMWlVsc+t42tm1HSEDmOs3A8fC9yyeisgLhsTVQokOZ0zw==", "license": "BSD-3-Clause", "dependencies": { - "@protobufjs/aspromise": "^1.1.1", - "@protobufjs/inquire": "^1.1.0" + "@protobufjs/aspromise": "^1.1.1" } }, "node_modules/@protobufjs/float": { @@ -331,12 +330,6 @@ "integrity": "sha512-Ddb+kVXlXst9d+R9PfTIxh1EdNkgoRe5tOX6t01f1lYWOvJnSPDBlG241QLzcyPdoNTsblLUdujGSE4RzrTZGQ==", "license": "BSD-3-Clause" }, - "node_modules/@protobufjs/inquire": { - "version": "1.1.0", - "resolved": "https://registry.npmjs.org/@protobufjs/inquire/-/inquire-1.1.0.tgz", - "integrity": "sha512-kdSefcPdruJiFMVSbn801t4vFK7KB/5gd2fYvrxhuJYg8ILrmn9SKSX2tZdV6V+ksulWqS7aXjBcRXl3wHoD9Q==", - "license": "BSD-3-Clause" - }, "node_modules/@protobufjs/path": { "version": "1.1.2", "resolved": "https://registry.npmjs.org/@protobufjs/path/-/path-1.1.2.tgz", @@ -350,9 +343,9 @@ "license": "BSD-3-Clause" }, "node_modules/@protobufjs/utf8": { - "version": "1.1.0", - "resolved": "https://registry.npmjs.org/@protobufjs/utf8/-/utf8-1.1.0.tgz", - "integrity": "sha512-Vvn3zZrhQZkkBE8LSuW3em98c0FwgO4nxzv6OdSxPKJIEKY2bGbHn+mhGIPerzI4twdxaP8/0+06HBpwf345Lw==", + "version": "1.1.2", + "resolved": "https://registry.npmjs.org/@protobufjs/utf8/-/utf8-1.1.2.tgz", + "integrity": "sha512-b1UQwcEZ4yCnMCD8DAL1VlbvBJE9/IX4FTIp7BG1xYpf29SLazLSrqUkj4w7Y5y7cCVP6E5tcqqcI0xemPkHug==", "license": "BSD-3-Clause" }, "node_modules/@scure/base": { @@ -3026,24 +3019,23 @@ } }, "node_modules/protobufjs": { - "version": "7.5.4", - "resolved": "https://registry.npmjs.org/protobufjs/-/protobufjs-7.5.4.tgz", - "integrity": "sha512-CvexbZtbov6jW2eXAvLukXjXUW1TzFaivC46BpWc/3BpcCysb5Vffu+B3XHMm8lVEuy2Mm4XGex8hBSg1yapPg==", + "version": "7.6.5", + "resolved": "https://registry.npmjs.org/protobufjs/-/protobufjs-7.6.5.tgz", + "integrity": "sha512-/FPD0nUc9jH6rfFjji9IBqOz4pcSE3CsT1m7Ep6Mdb0LxSUMj8hgl6GomOvZzpNpAqqGaXA0P3VSrZLFzIhQrw==", "hasInstallScript": true, "license": "BSD-3-Clause", "dependencies": { "@protobufjs/aspromise": "^1.1.2", "@protobufjs/base64": "^1.1.2", - "@protobufjs/codegen": "^2.0.4", - "@protobufjs/eventemitter": "^1.1.0", - "@protobufjs/fetch": "^1.1.0", + "@protobufjs/codegen": "^2.0.5", + "@protobufjs/eventemitter": "^1.1.1", + "@protobufjs/fetch": "^1.1.1", "@protobufjs/float": "^1.0.2", - "@protobufjs/inquire": "^1.1.0", "@protobufjs/path": "^1.1.2", "@protobufjs/pool": "^1.1.0", - "@protobufjs/utf8": "^1.1.0", + "@protobufjs/utf8": "^1.1.1", "@types/node": ">=13.7.0", - "long": "^5.0.0" + "long": "^5.3.2" }, "engines": { "node": ">=12.0.0"