diff --git a/barretenberg/cpp/src/barretenberg/crypto/merkle_tree/lmdb_store/lmdb_tree_store.cpp b/barretenberg/cpp/src/barretenberg/crypto/merkle_tree/lmdb_store/lmdb_tree_store.cpp index 6127f8040630..34bb07c30b75 100644 --- a/barretenberg/cpp/src/barretenberg/crypto/merkle_tree/lmdb_store/lmdb_tree_store.cpp +++ b/barretenberg/cpp/src/barretenberg/crypto/merkle_tree/lmdb_store/lmdb_tree_store.cpp @@ -54,8 +54,9 @@ int index_key_cmp(const MDB_val* a, const MDB_val* b) return value_cmp(a, b); } -LMDBTreeStore::LMDBTreeStore(std::string directory, std::string name, uint64_t mapSizeKb, uint64_t maxNumReaders) - : LMDBStoreBase(directory, mapSizeKb, maxNumReaders, 5) +LMDBTreeStore::LMDBTreeStore( + std::string directory, std::string name, uint64_t mapSizeKb, uint64_t maxNumReaders, bool ephemeral) + : LMDBStoreBase(directory, mapSizeKb, maxNumReaders, 5, ephemeral) , _name(std::move(name)) { diff --git a/barretenberg/cpp/src/barretenberg/crypto/merkle_tree/lmdb_store/lmdb_tree_store.hpp b/barretenberg/cpp/src/barretenberg/crypto/merkle_tree/lmdb_store/lmdb_tree_store.hpp index 30f68061b706..c916924f4e69 100644 --- a/barretenberg/cpp/src/barretenberg/crypto/merkle_tree/lmdb_store/lmdb_tree_store.hpp +++ b/barretenberg/cpp/src/barretenberg/crypto/merkle_tree/lmdb_store/lmdb_tree_store.hpp @@ -159,7 +159,8 @@ class LMDBTreeStore : public LMDBStoreBase { using SharedPtr = std::shared_ptr; using ReadTransaction = LMDBReadTransaction; using WriteTransaction = LMDBWriteTransaction; - LMDBTreeStore(std::string directory, std::string name, uint64_t mapSizeKb, uint64_t maxNumReaders); + LMDBTreeStore( + std::string directory, std::string name, uint64_t mapSizeKb, uint64_t maxNumReaders, bool ephemeral = false); LMDBTreeStore(const LMDBTreeStore& other) = delete; LMDBTreeStore(LMDBTreeStore&& other) = delete; LMDBTreeStore& operator=(const LMDBTreeStore& other) = delete; diff --git a/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_environment.cpp b/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_environment.cpp index 99571f661c49..ab712c7fa2b5 100644 --- a/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_environment.cpp +++ b/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_environment.cpp @@ -8,10 +8,8 @@ namespace bb::lmdblib { -LMDBEnvironment::LMDBEnvironment(const std::string& directory, - uint64_t mapSizeKB, - uint32_t maxNumDBs, - uint32_t maxNumReaders) +LMDBEnvironment::LMDBEnvironment( + const std::string& directory, uint64_t mapSizeKB, uint32_t maxNumDBs, uint32_t maxNumReaders, bool ephemeral) : _id(0) , _directory(directory) , _readGuard(maxNumReaders) @@ -21,6 +19,9 @@ LMDBEnvironment::LMDBEnvironment(const std::string& directory, uint64_t kb = 1024; uint64_t totalMapSize = kb * mapSizeKB; uint32_t flags = MDB_NOTLS; + if (ephemeral) { + flags |= MDB_NOSYNC | MDB_NOMETASYNC; + } try { call_lmdb_func("mdb_env_set_mapsize", mdb_env_set_mapsize, _mdbEnv, static_cast(totalMapSize)); call_lmdb_func("mdb_env_set_maxdbs", mdb_env_set_maxdbs, _mdbEnv, static_cast(maxNumDBs)); diff --git a/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_environment.hpp b/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_environment.hpp index c87054028ac4..c463ca74675c 100644 --- a/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_environment.hpp +++ b/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_environment.hpp @@ -25,8 +25,19 @@ class LMDBEnvironment { * @param mapSizeKb The maximum size of the database, can be increased from a previously used value * @param maxNumDbs The maximum number of databases that can be created withn this environment * @param maxNumReaders The maximum number of concurrent read transactions permitted. + * @param ephemeral When true, opens the env with `MDB_NOSYNC | MDB_NOMETASYNC`. Commits + * return as soon as the dirty pages are queued; the kernel flushes them + * lazily and never blocks the commit. Files stay sparse (we deliberately + * avoid `MDB_WRITEMAP`, which would eagerly allocate the full map size + * on disk). Intended for short-lived ephemeral world states (e.g. TXE + * test sessions) that discard the directory on close; never use for a + * real node — a crash mid-write yields an unrecoverable env. */ - LMDBEnvironment(const std::string& directory, uint64_t mapSizeKb, uint32_t maxNumDBs, uint32_t maxNumReaders); + LMDBEnvironment(const std::string& directory, + uint64_t mapSizeKb, + uint32_t maxNumDBs, + uint32_t maxNumReaders, + bool ephemeral = false); LMDBEnvironment(const LMDBEnvironment& other) = delete; LMDBEnvironment(LMDBEnvironment&& other) = delete; LMDBEnvironment& operator=(const LMDBEnvironment& other) = delete; diff --git a/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_store.cpp b/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_store.cpp index ef886c47e147..15165a93d910 100644 --- a/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_store.cpp +++ b/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_store.cpp @@ -14,8 +14,8 @@ #include namespace bb::lmdblib { -LMDBStore::LMDBStore(std::string directory, uint64_t mapSizeKb, uint64_t maxNumReaders, uint64_t maxDbs) - : LMDBStoreBase(std::move(directory), mapSizeKb, maxNumReaders, maxDbs) +LMDBStore::LMDBStore(std::string directory, uint64_t mapSizeKb, uint64_t maxNumReaders, uint64_t maxDbs, bool ephemeral) + : LMDBStoreBase(std::move(directory), mapSizeKb, maxNumReaders, maxDbs, ephemeral) {} void LMDBStore::open_database(const std::string& name, bool duplicateKeysPermitted) diff --git a/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_store.hpp b/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_store.hpp index ffcf8d8cdc33..299dc9f65c12 100644 --- a/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_store.hpp +++ b/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_store.hpp @@ -35,7 +35,8 @@ class LMDBStore : public LMDBStoreBase { std::string name; }; - LMDBStore(std::string directory, uint64_t mapSizeKb, uint64_t maxNumReaders, uint64_t maxDbs); + LMDBStore( + std::string directory, uint64_t mapSizeKb, uint64_t maxNumReaders, uint64_t maxDbs, bool ephemeral = false); LMDBStore(const LMDBStore& other) = delete; LMDBStore(LMDBStore&& other) = delete; LMDBStore& operator=(const LMDBStore& other) = delete; diff --git a/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_store_base.cpp b/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_store_base.cpp index 9bc2bab4cf2c..459510da23bb 100644 --- a/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_store_base.cpp +++ b/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_store_base.cpp @@ -1,9 +1,10 @@ #include "barretenberg/lmdblib/lmdb_store_base.hpp" namespace bb::lmdblib { -LMDBStoreBase::LMDBStoreBase(std::string directory, uint64_t mapSizeKb, uint64_t maxNumReaders, uint64_t maxDbs) +LMDBStoreBase::LMDBStoreBase( + std::string directory, uint64_t mapSizeKb, uint64_t maxNumReaders, uint64_t maxDbs, bool ephemeral) : _dbDirectory(std::move(directory)) - , _environment((std::make_shared(_dbDirectory, mapSizeKb, maxDbs, maxNumReaders))) + , _environment(std::make_shared(_dbDirectory, mapSizeKb, maxDbs, maxNumReaders, ephemeral)) {} LMDBStoreBase::~LMDBStoreBase() = default; LMDBStoreBase::ReadTransaction::Ptr LMDBStoreBase::create_read_transaction() const diff --git a/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_store_base.hpp b/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_store_base.hpp index f9a40cd34277..fd03e29c4ea4 100644 --- a/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_store_base.hpp +++ b/barretenberg/cpp/src/barretenberg/lmdblib/lmdb_store_base.hpp @@ -11,7 +11,8 @@ class LMDBStoreBase { using ReadTransaction = LMDBReadTransaction; using WriteTransaction = LMDBWriteTransaction; using DBCreationTransaction = LMDBDatabaseCreationTransaction; - LMDBStoreBase(std::string directory, uint64_t mapSizeKb, uint64_t maxNumReaders, uint64_t maxDbs); + LMDBStoreBase( + std::string directory, uint64_t mapSizeKb, uint64_t maxNumReaders, uint64_t maxDbs, bool ephemeral = false); LMDBStoreBase(const LMDBStoreBase& other) = delete; LMDBStoreBase& operator=(const LMDBStoreBase& other) = delete; LMDBStoreBase(LMDBStoreBase&& other) noexcept = default; diff --git a/barretenberg/cpp/src/barretenberg/nodejs_module/lmdb_store/lmdb_store_wrapper.cpp b/barretenberg/cpp/src/barretenberg/nodejs_module/lmdb_store/lmdb_store_wrapper.cpp index 4aea32c58f58..b891f0e3f115 100644 --- a/barretenberg/cpp/src/barretenberg/nodejs_module/lmdb_store/lmdb_store_wrapper.cpp +++ b/barretenberg/cpp/src/barretenberg/nodejs_module/lmdb_store/lmdb_store_wrapper.cpp @@ -56,7 +56,21 @@ LMDBStoreWrapper::LMDBStoreWrapper(const Napi::CallbackInfo& info) } } - _store = std::make_unique(data_dir, map_size, max_readers, 2); + // `ephemeral` opens the LMDB env with `MDB_NOSYNC | MDB_NOMETASYNC`, so commits return + // without waiting for fsync. The on-disk file is unrecoverable after a crash but stays + // sparse on every byte LMDB doesn't touch; the trade-off is appropriate for tmp stores + // that get cleaned up on close (see `openTmpStore`). + size_t ephemeral_index = 3; + bool ephemeral = false; + if (info.Length() > ephemeral_index) { + if (info[ephemeral_index].IsBoolean()) { + ephemeral = info[ephemeral_index].As().Value(); + } else if (!info[ephemeral_index].IsUndefined()) { + throw Napi::TypeError::New(env, "The ephemeral flag must be a boolean"); + } + } + + _store = std::make_unique(data_dir, map_size, max_readers, 2, ephemeral); _msg_processor.register_handler(LMDBStoreMessageType::OPEN_DATABASE, this, &LMDBStoreWrapper::open_database); diff --git a/barretenberg/cpp/src/barretenberg/nodejs_module/world_state/world_state.cpp b/barretenberg/cpp/src/barretenberg/nodejs_module/world_state/world_state.cpp index 848e3b234096..cf55d7c5f6ec 100644 --- a/barretenberg/cpp/src/barretenberg/nodejs_module/world_state/world_state.cpp +++ b/barretenberg/cpp/src/barretenberg/nodejs_module/world_state/world_state.cpp @@ -169,6 +169,18 @@ WorldStateWrapper::WorldStateWrapper(const Napi::CallbackInfo& info) thread_pool_size = info[thread_pool_size_index].As().Uint32Value(); } + // `ephemeral` opens each underlying LMDB env with `MDB_NOSYNC | MDB_NOMETASYNC` — + // commits never block on fsync, files stay sparse, and a crash mid-write yields an + // unrecoverable env. Intended for throwaway scratch state (TXE test sessions). + bool ephemeral = false; + size_t ephemeral_index = 8; + if (info.Length() > ephemeral_index) { + if (!info[ephemeral_index].IsBoolean()) { + throw Napi::TypeError::New(env, "Ephemeral flag must be a boolean"); + } + ephemeral = info[ephemeral_index].As().Value(); + } + _ws = std::make_unique(thread_pool_size, data_dir, map_size, @@ -176,7 +188,8 @@ WorldStateWrapper::WorldStateWrapper(const Napi::CallbackInfo& info) tree_prefill, prefilled_public_data, initial_header_generator_point, - genesis_timestamp); + genesis_timestamp, + ephemeral); _dispatcher.register_target( WorldStateMessageType::GET_TREE_INFO, diff --git a/barretenberg/cpp/src/barretenberg/world_state/world_state.cpp b/barretenberg/cpp/src/barretenberg/world_state/world_state.cpp index 49cf9229870d..883fefbd0e9b 100644 --- a/barretenberg/cpp/src/barretenberg/world_state/world_state.cpp +++ b/barretenberg/cpp/src/barretenberg/world_state/world_state.cpp @@ -40,7 +40,8 @@ WorldState::WorldState(uint64_t thread_pool_size, const std::unordered_map& tree_prefill, const std::vector& prefilled_public_data, uint32_t initial_header_generator_point, - uint64_t genesis_timestamp) + uint64_t genesis_timestamp, + bool ephemeral) : _workers(std::make_shared(thread_pool_size)) , _tree_heights(tree_heights) , _initial_tree_size(tree_prefill) @@ -50,7 +51,7 @@ WorldState::WorldState(uint64_t thread_pool_size, { // We set the max readers to be high, at least the number of given threads or the default if higher uint64_t maxReaders = std::max(thread_pool_size, DEFAULT_MIN_NUMBER_OF_READERS); - create_canonical_fork(data_dir, map_size, prefilled_public_data, maxReaders); + create_canonical_fork(data_dir, map_size, prefilled_public_data, maxReaders, ephemeral); try { attempt_tree_resync(); } catch (std::exception& e) { @@ -64,7 +65,8 @@ WorldState::WorldState(uint64_t thread_pool_size, const std::unordered_map& tree_heights, const std::unordered_map& tree_prefill, uint32_t initial_header_generator_point, - uint64_t genesis_timestamp) + uint64_t genesis_timestamp, + bool ephemeral) : WorldState::WorldState(thread_pool_size, data_dir, map_size, @@ -72,7 +74,8 @@ WorldState::WorldState(uint64_t thread_pool_size, tree_prefill, std::vector(), initial_header_generator_point, - genesis_timestamp) + genesis_timestamp, + ephemeral) {} WorldState::WorldState(uint64_t thread_pool_size, @@ -82,7 +85,8 @@ WorldState::WorldState(uint64_t thread_pool_size, const std::unordered_map& tree_prefill, const std::vector& prefilled_public_data, uint32_t initial_header_generator_point, - uint64_t genesis_timestamp) + uint64_t genesis_timestamp, + bool ephemeral) : WorldState(thread_pool_size, data_dir, { @@ -96,7 +100,8 @@ WorldState::WorldState(uint64_t thread_pool_size, tree_prefill, prefilled_public_data, initial_header_generator_point, - genesis_timestamp) + genesis_timestamp, + ephemeral) {} WorldState::WorldState(uint64_t thread_pool_size, @@ -105,7 +110,8 @@ WorldState::WorldState(uint64_t thread_pool_size, const std::unordered_map& tree_heights, const std::unordered_map& tree_prefill, uint32_t initial_header_generator_point, - uint64_t genesis_timestamp) + uint64_t genesis_timestamp, + bool ephemeral) : WorldState(thread_pool_size, data_dir, map_size, @@ -113,13 +119,15 @@ WorldState::WorldState(uint64_t thread_pool_size, tree_prefill, std::vector(), initial_header_generator_point, - genesis_timestamp) + genesis_timestamp, + ephemeral) {} void WorldState::create_canonical_fork(const std::string& dataDir, const std::unordered_map& dbSize, const std::vector& prefilled_public_data, - uint64_t maxReaders) + uint64_t maxReaders, + bool ephemeral) { // create the underlying stores auto createStore = [&](MerkleTreeId id) { @@ -127,7 +135,7 @@ void WorldState::create_canonical_fork(const std::string& dataDir, std::filesystem::path directory = dataDir; directory /= name; std::filesystem::create_directories(directory); - return std::make_shared(directory, name, dbSize.at(id), maxReaders); + return std::make_shared(directory, name, dbSize.at(id), maxReaders, ephemeral); }; _persistentStores = std::make_unique(createStore(MerkleTreeId::NULLIFIER_TREE), createStore(MerkleTreeId::PUBLIC_DATA_TREE), diff --git a/barretenberg/cpp/src/barretenberg/world_state/world_state.hpp b/barretenberg/cpp/src/barretenberg/world_state/world_state.hpp index cdd825801778..2e493e2efbdf 100644 --- a/barretenberg/cpp/src/barretenberg/world_state/world_state.hpp +++ b/barretenberg/cpp/src/barretenberg/world_state/world_state.hpp @@ -64,7 +64,8 @@ class WorldState { const std::unordered_map& tree_heights, const std::unordered_map& tree_prefill, uint32_t initial_header_generator_point, - uint64_t genesis_timestamp = 0); + uint64_t genesis_timestamp = 0, + bool ephemeral = false); WorldState(uint64_t thread_pool_size, const std::string& data_dir, @@ -72,7 +73,8 @@ class WorldState { const std::unordered_map& tree_heights, const std::unordered_map& tree_prefill, uint32_t initial_header_generator_point, - uint64_t genesis_timestamp = 0); + uint64_t genesis_timestamp = 0, + bool ephemeral = false); WorldState(uint64_t thread_pool_size, const std::string& data_dir, @@ -81,8 +83,16 @@ class WorldState { const std::unordered_map& tree_prefill, const std::vector& prefilled_public_data, uint32_t initial_header_generator_point, - uint64_t genesis_timestamp = 0); + uint64_t genesis_timestamp = 0, + bool ephemeral = false); + /** + * @param ephemeral When true, every underlying LMDB env opens with `MDB_NOSYNC | + * MDB_NOMETASYNC`. Commits return without waiting for fsync; the kernel + * flushes lazily, files stay sparse. Intended for throwaway scratch + * state. A crash mid-write yields an unrecoverable + * env, so never enable on a node that needs to survive a restart. + */ WorldState(uint64_t thread_pool_size, const std::string& data_dir, const std::unordered_map& map_size, @@ -90,7 +100,8 @@ class WorldState { const std::unordered_map& tree_prefill, const std::vector& prefilled_public_data, uint32_t initial_header_generator_point, - uint64_t genesis_timestamp = 0); + uint64_t genesis_timestamp = 0, + bool ephemeral = false); /** * @brief Copies all underlying LMDB stores to the target directory while acquiring a write lock @@ -313,7 +324,8 @@ class WorldState { void create_canonical_fork(const std::string& dataDir, const std::unordered_map& dbSize, const std::vector& prefilled_public_data, - uint64_t maxReaders); + uint64_t maxReaders, + bool ephemeral); Fork::SharedPtr retrieve_fork(const uint64_t& forkId) const; Fork::SharedPtr create_new_fork(const block_number_t& blockNumber); diff --git a/yarn-project/aztec/scripts/aztec.sh b/yarn-project/aztec/scripts/aztec.sh index 67eddafe9516..0477c3923f94 100755 --- a/yarn-project/aztec/scripts/aztec.sh +++ b/yarn-project/aztec/scripts/aztec.sh @@ -33,7 +33,7 @@ case $cmd in fi while ! nc -z 127.0.0.1 8081 &>/dev/null; do sleep 0.2; done export NARGO_FOREIGN_CALL_TIMEOUT=300000 - nargo test --silence-warnings --oracle-resolver http://127.0.0.1:8081 --test-threads 16 "$@" + nargo test --silence-warnings --oracle-resolver http://127.0.0.1:8081 "$@" ;; start) if [ "${1:-}" == "--local-network" ]; then diff --git a/yarn-project/aztec/src/cli/cmds/start_txe.ts b/yarn-project/aztec/src/cli/cmds/start_txe.ts index 20ae1b56a7ee..6266b0691810 100644 --- a/yarn-project/aztec/src/cli/cmds/start_txe.ts +++ b/yarn-project/aztec/src/cli/cmds/start_txe.ts @@ -1,11 +1,11 @@ import { startHttpRpcServer } from '@aztec/foundation/json-rpc/server'; import type { Logger } from '@aztec/foundation/log'; -import { createTXERpcServer } from '@aztec/txe'; +import { createTXERpcServer } from '@aztec/txe/server'; export async function startTXE(options: any, signalHandlers: Array<() => Promise>, debugLogger: Logger) { debugLogger.info(`Setting up TXE...`); - const txeServer = createTXERpcServer(debugLogger); + const txeServer = await createTXERpcServer(debugLogger); const httpServer = await startHttpRpcServer(txeServer, { port: options.port, timeoutMs: 1e3 * 60 * 5, diff --git a/yarn-project/bootstrap.sh b/yarn-project/bootstrap.sh index 46c0c712bdac..5f0d8957c49a 100755 --- a/yarn-project/bootstrap.sh +++ b/yarn-project/bootstrap.sh @@ -143,6 +143,9 @@ function compile_all { get_projects | compile_project + cd txe && yarn build + cd .. + # Run oracle version checks after compilation cd pxe && yarn check_oracle_version cd .. diff --git a/yarn-project/kv-store/src/lmdb-v2/factory.ts b/yarn-project/kv-store/src/lmdb-v2/factory.ts index eb668537fe86..407caa4f8640 100644 --- a/yarn-project/kv-store/src/lmdb-v2/factory.ts +++ b/yarn-project/kv-store/src/lmdb-v2/factory.ts @@ -3,7 +3,7 @@ import { type LoggerBindings, createLogger } from '@aztec/foundation/log'; import { DatabaseVersionManager, type SchemaVersionMismatchPolicy } from '@aztec/stdlib/database-version/manager'; import type { DataStoreConfig } from '@aztec/stdlib/kv-store'; -import { mkdir, mkdtemp, rm } from 'fs/promises'; +import { copyFile, mkdir, mkdtemp, rm } from 'fs/promises'; import { tmpdir } from 'os'; import { join } from 'path'; @@ -57,9 +57,14 @@ export async function createStore( return store; } +/** + * Open a persistent on-disk store rooted in OS tmpdir. Caller chooses whether to + * auto-remove the directory on close. Uses standard durable LMDB flags (fsync on + * every commit) — for full in-memory / NOSYNC semantics use `openEphemeralStore`. + */ export async function openTmpStore( name: string, - ephemeral: boolean = true, + cleanupTmpDir: boolean = true, dbMapSizeKb = 10 * 1_024 * 1_024, // 10GB maxReaders = MAX_READERS, bindings?: LoggerBindings, @@ -70,7 +75,7 @@ export async function openTmpStore( // pass a cleanup callback because process.on('beforeExit', cleanup) does not work under Jest const cleanup = async () => { - if (ephemeral) { + if (cleanupTmpDir) { try { await rm(dataDir, { recursive: true, force: true, maxRetries: 3 }); log.debug(`Deleted temporary data store: ${dataDir}`); @@ -82,9 +87,35 @@ export async function openTmpStore( } }; - // For temporary stores, we don't need to worry about versioning - // as they are ephemeral and get cleaned up after use - return AztecLMDBStoreV2.new(dataDir, dbMapSizeKb, maxReaders, cleanup, bindings); + return AztecLMDBStoreV2.new(dataDir, dbMapSizeKb, maxReaders, cleanup, bindings, false); +} + +/** + * Open a fully-ephemeral store: a fresh tmpdir backing file, LMDB opened with + * MDB_NOSYNC | MDB_NOMETASYNC (no fsync, kernel-only writeback), and the directory + * unconditionally removed on close. Intended for tests and worker-local scratch + * stores where durability is explicitly not desired. + */ +export async function openEphemeralStore( + name: string, + dbMapSizeKb = 10 * 1_024 * 1_024, // 10GB + maxReaders = MAX_READERS, + bindings?: LoggerBindings, +): Promise { + const log = createLogger('kv-store:lmdb-v2:' + name, bindings); + const dataDir = await mkdtemp(join(tmpdir(), name + '-')); + log.debug(`Created ephemeral data store at: ${dataDir} with size: ${dbMapSizeKb} KB (LMDB v2)`); + + const cleanup = async () => { + try { + await rm(dataDir, { recursive: true, force: true, maxRetries: 3 }); + log.debug(`Deleted ephemeral data store: ${dataDir}`); + } catch (err) { + log.warn(`Failed to delete ephemeral data directory (LMDB v2) ${dataDir}: ${err}`); + } + }; + + return AztecLMDBStoreV2.new(dataDir, dbMapSizeKb, maxReaders, cleanup, bindings, true); } export async function openStoreAt( @@ -98,6 +129,36 @@ export async function openStoreAt( return await AztecLMDBStoreV2.new(dataDir, dbMapSizeKb, maxReaders, undefined, bindings); } +/** + * Open a fully-ephemeral store seeded from an existing `data.mdb` file. Creates a fresh tmpdir, + * copies `srcDataMdbPath` into it as `data.mdb`, opens with MDB_NOSYNC | MDB_NOMETASYNC, and + * removes the tmpdir on close. Intended for worker clones of a main-thread-built store — the + * source file is never written to. + */ +export async function cloneEphemeralStoreFrom( + srcDataMdbPath: string, + name: string, + dbMapSizeKb = 10 * 1_024 * 1_024, // 10GB + maxReaders = MAX_READERS, + bindings?: LoggerBindings, +): Promise { + const log = createLogger('kv-store:lmdb-v2:' + name, bindings); + const dataDir = await mkdtemp(join(tmpdir(), name + '-')); + await copyFile(srcDataMdbPath, join(dataDir, 'data.mdb')); + log.debug(`Cloned ephemeral data store at: ${dataDir} from ${srcDataMdbPath} (LMDB v2)`); + + const cleanup = async () => { + try { + await rm(dataDir, { recursive: true, force: true, maxRetries: 3 }); + log.debug(`Deleted ephemeral data store: ${dataDir}`); + } catch (err) { + log.warn(`Failed to delete ephemeral data directory (LMDB v2) ${dataDir}: ${err}`); + } + }; + + return AztecLMDBStoreV2.new(dataDir, dbMapSizeKb, maxReaders, cleanup, bindings, true); +} + export async function openVersionedStoreAt( dataDirectory: string, schemaVersion: number, diff --git a/yarn-project/kv-store/src/lmdb-v2/store.ts b/yarn-project/kv-store/src/lmdb-v2/store.ts index d7265a6f422d..51c1011ddc1e 100644 --- a/yarn-project/kv-store/src/lmdb-v2/store.ts +++ b/yarn-project/kv-store/src/lmdb-v2/store.ts @@ -43,9 +43,10 @@ export class AztecLMDBStoreV2 implements AztecAsyncKVStore, LMDBMessageChannel { maxReaders: number, private log: Logger, private cleanup?: () => Promise, + ephemeral: boolean = false, ) { this.log.info(`Starting data store with maxReaders ${maxReaders}`); - this.channel = new MsgpackChannel(new NativeLMDBStore(dataDir, mapSize, maxReaders)); + this.channel = new MsgpackChannel(new NativeLMDBStore(dataDir, mapSize, maxReaders, ephemeral)); // leave one reader to always be available for regular, atomic, reads this.availableCursors = new Semaphore(maxReaders - 1); } @@ -76,9 +77,10 @@ export class AztecLMDBStoreV2 implements AztecAsyncKVStore, LMDBMessageChannel { maxReaders: number = 16, cleanup?: () => Promise, bindings?: LoggerBindings, + ephemeral: boolean = false, ) { const log = createLogger('kv-store:lmdb-v2', bindings); - const db = new AztecLMDBStoreV2(dataDir, dbMapSizeKb, maxReaders, log, cleanup); + const db = new AztecLMDBStoreV2(dataDir, dbMapSizeKb, maxReaders, log, cleanup, ephemeral); await db.start(); return db; } diff --git a/yarn-project/txe/esbuild.config.mjs b/yarn-project/txe/esbuild.config.mjs new file mode 100644 index 000000000000..46be5f765882 --- /dev/null +++ b/yarn-project/txe/esbuild.config.mjs @@ -0,0 +1,243 @@ +// Bundles the TXE entry points (bin/index.ts + worker.ts + rpc_server.ts) into single files so +// each Node isolate — main thread and worker_threads alike — compiles one file at startup instead +// of resolving the full TXE → simulator → world-state → bb-prover dependency tree from scratch. +// This is the dominant cost of cold-starting TXE. +// +// Native modules (LMDB, @aztec/native, msgpackr, snappy, etc.) and packages that load their own +// .wasm assets at runtime (@aztec/noir-acvm_js, @aztec/noir-noirc_abi, @aztec/bb.js) stay +// external — esbuild cannot bundle .node binaries, and bundling the WASM-loader JS would break +// relative paths to the .wasm files. +// +import { build } from 'esbuild'; +import { writeFile } from 'node:fs/promises'; + +import { externalNativePlugin } from './esbuild/plugins/external_native.mjs'; +import { protocolContractsEagerToLazyPlugin } from './esbuild/plugins/protocol_contracts_eager_to_lazy.mjs'; +import { redirectsPlugin } from './esbuild/plugins/redirects.mjs'; +import { enforceSizeLimits } from './esbuild/plugins/size_guard.mjs'; +import { stripArtifactDebugPlugin } from './esbuild/plugins/strip_artifact_debug.mjs'; + +// `@aztec/*` packages that AztecNodeService imports transitively but the TXE worker never +// actually constructs or executes — sequencer/validator/prover/etc. are passed in as `undefined`, +// so the imported symbols are only used inside AztecNodeService methods TXE never calls. The +// Proxy-based `empty_stub.cjs` throws on access, surfacing any false assumption loudly. +// +// NOTE: `@aztec/archiver` is intentionally NOT stubbed — TXEArchiver extends ArchiverDataSourceBase +// and uses createArchiverDataStores at runtime. +// NOTE: `@aztec/blob-lib` is NOT stubbed — stdlib's tx_effect calls getNumTxBlobFields at runtime +// when serializing transactions. +const fullyStubbedAztecPackages = [ + '@aztec/p2p', + '@aztec/sequencer-client', + '@aztec/validator-client', + '@aztec/prover-client', + '@aztec/prover-node', + '@aztec/slasher', + '@aztec/epoch-cache', + '@aztec/blob-client', + '@aztec/node-keystore', + '@aztec/node-lib', +]; + +// Each row redirects a specifier match to a stub file under `esbuild/stubs/`. The local stub +// re-exports only the surface the bundled code reaches at runtime, or throws on call for code +// paths TXE never enters (see `throwTrap` in foundation/error). For relative-path filters, +// `importerContains` gates by importer path so we don't hijack same-named files in other +// packages. +const redirects = [ + // Whole-package redirects (bare specifier + every subpath → shared Proxy stub). + { + filter: new RegExp(`^(${fullyStubbedAztecPackages.map(p => p.replace(/[/-]/g, '[\\/\\-]')).join('|')})(/.*)?$`), + stub: 'empty_stub.cjs', + }, + + // Bare-specifier redirects: each package has its own bespoke stub re-exporting just what the + // bundled code uses. The real packages either drag in heavy transitive deps the TXE worker + // never executes (telemetry HTTP wrappers, archiver factory + L1 retrieval, world-state + // synchronizer, bb prover/verifier stack) or eagerly precompute large lookup tables at import + // (noble curve field-tower tables). + { filter: /^@aztec\/telemetry-client$/, stub: 'telemetry_stub.ts' }, + { filter: /^@aztec\/protocol-contracts\/providers\/bundle$/, stub: 'protocol_contracts_bundle_stub.ts' }, + { filter: /^@aztec\/archiver$/, stub: 'archiver_stub.ts' }, + { filter: /^@aztec\/bb-prover$/, stub: 'bb_prover_stub.ts' }, + { filter: /^@aztec\/bb-prover\/test$/, stub: 'bb_prover_test_stub.ts' }, + { filter: /^@aztec\/world-state$/, stub: 'world_state_stub.ts' }, + { filter: /^@noble\/curves\/secp256k1$/, stub: 'noble_secp256k1_stub.ts' }, + { filter: /^@noble\/curves\/bls12-381$/, stub: 'noble_bls12_stub.ts' }, + + // Relative-path redirects with importer gating: a specific file inside a package would + // otherwise pull in heavy unused code, but we can't wholesale stub the package (it has live + // siblings TXE still needs). + // + // - aztec-node/.../node_metrics.js: constructed by AztecNodeService but never observed + // because TXE's meter is the Noop client. + // - aztec-node/.../sentinel/{factory,sentinel}.js: only reached via `AztecNodeService.start()`, + // which TXE never invokes. + // - telemetry-client/.../start.js: contains `await import('./otel.js')` which forces emission + // of a large OpenTelemetry SDK chunk even though TXE never starts telemetry. + // - zod/.../locales/index.js: barrel re-exports every i18n bundle; TXE only needs `en`. + { + filter: /^\.\/node_metrics\.(js|ts)$/, + importerContains: '/aztec-node/', + stub: 'aztec_node_node_metrics_stub.ts', + }, + { + filter: /^\.\.\/sentinel\/factory\.(js|ts)$/, + importerContains: '/aztec-node/', + stub: 'aztec_node_sentinel_factory_stub.ts', + }, + { + filter: /^\.\.\/sentinel\/sentinel\.(js|ts)$/, + importerContains: '/aztec-node/', + stub: 'aztec_node_sentinel_stub.ts', + }, + // start.js arrives as a resolved absolute path from outside the package… + { filter: /telemetry-client\/(dest|src)\/start\.(js|ts)$/, stub: 'telemetry_start_stub.ts' }, + // …and as a relative `./start.js` from siblings inside the package. + { + filter: /^\.\.?\/start\.js$/, + importerContains: '/telemetry-client/', + stub: 'telemetry_start_stub.ts', + }, + { filter: /locales\/index\.js$/, importerContains: '/zod/', stub: 'zod_locales_stub.ts' }, + // Resolved-path form of providers/bundle.js (the bare-specifier form is matched above). + { + filter: /protocol-contracts\/(dest|src)\/provider\/bundle\.(js|ts)$/, + stub: 'protocol_contracts_bundle_stub.ts', + }, +]; + +const entryPoints = { + // src/bin/index.ts → dest/bin/index.js (overwrites the tsc-emitted file). + 'bin/index': 'src/bin/index.ts', + // src/worker.ts → dest/worker.bundle.js (the file the pool spawns). + 'worker.bundle': 'src/worker.ts', + // src/rpc_server.ts → dest/server.bundle.js (the entry the parent `@aztec/aztec` + // CLI imports via `@aztec/txe/server`). + 'server.bundle': 'src/rpc_server.ts', +}; + +const start = Date.now(); +const result = await build({ + entryPoints, + outdir: 'dest', + entryNames: '[dir]/[name]', + bundle: true, + platform: 'node', + format: 'esm', + // Splitting lets `LazyProtocolContractsProvider`'s per-contract `await import(...)` calls + // resolve into separate chunks loaded only when warmUp() runs. Without splitting, esbuild + // inlines dynamic imports back into the parent bundle and the artifacts are eager again. + splitting: true, + target: 'node20', + external: [ + // `pako` is used only by `parseDebugSymbols` in stdlib/abi/abi.ts. We strip `debug_symbols` + // at bundle time, so `getFunctionDebugMetadata` always short-circuits before reaching + // parseDebugSymbols — pako never executes. Externalizing keeps the inflate impl out of the + // eager startup chunk; if it ever IS reached, Node will resolve it from node_modules. + 'pako', + // Native-wrapper JS packages. Bundling these would either pull in their .node files (which + // can't be loaded by esbuild) or strand them from their per-arch native dependencies. + 'lmdb', + 'msgpackr', + 'msgpackr-extract', + 'snappy', + 'node-eth-kzg', + '@aztec/native', + // WASM-loading packages. + '@aztec/noir-acvm_js', + '@aztec/noir-noirc_abi', + // bb.js loads barretenberg-threads.wasm.gz via a path computed relative to its own module + // location — bundling moves the JS but not the .wasm.gz, breaking the load. + '@aztec/bb.js', + // Artifacts the TXE worker never executes (rollup, parity, kernel circuits, VKs). They ship + // as large JSON files and get re-exported from a barrel without proper `sideEffects: false`, + // so esbuild can't tree-shake them. Externalizing the whole package keeps them off the + // worker's startup parse path. + '@aztec/noir-protocol-circuits-types', + '@aztec/noir-protocol-circuits-types/client', + '@aztec/noir-protocol-circuits-types/server', + '@aztec/noir-protocol-circuits-types/vks', + // pino spawns its own worker_threads for transports and detects its runtime format + // dynamically; bundling it triggers ERR_AMBIGUOUS_MODULE_SYNTAX. thread-stream and + // sonic-boom are pino's siblings and have the same issue. + 'pino', + 'pino-pretty', + 'pino-abstract-transport', + 'thread-stream', + 'sonic-boom', + // bcrypto detects native vs js bindings at load time the same way pino does; bundling + // triggers ERR_AMBIGUOUS_MODULE_SYNTAX. + 'bcrypto', + // TXE never talks to L1, so the whole ethereum/viem stack is dead weight in the worker. + // Keeping it external avoids bundling forwarder_proxy (whose CLI guard executes when + // bundled), viem, abitype, and the rest of the L1 client surface. + '@aztec/ethereum', + '@aztec/l1-artifacts', + // Same reasoning: stdlib/p2p/signature_utils, archiver/l1, and aztec-node/config all + // statically import `viem` for L1 RPC + signing primitives. None of those code paths execute + // in TXE, so externalize the whole tree. + 'viem', + 'abitype', + 'ox', + // TXE never snapshots/uploads world state, so the file-store backend (which drags in the + // AWS S3 SDK, Google Cloud Storage SDK, mime-db, pako, etc.) is dead weight. + '@aztec/stdlib/file-store', + '@aztec/stdlib/snapshots', + ], + plugins: [ + externalNativePlugin, + redirectsPlugin(redirects), + stripArtifactDebugPlugin, + protocolContractsEagerToLazyPlugin, + ], + // esbuild's `__require` helper looks for a top-level `require` binding in the bundle and falls + // back to throwing `Dynamic require of "X" is not supported` if absent. We must provide one for + // bundled CJS deps (e.g. pino's `require('node:os')`). Imported under a renamed alias so we do + // not collide with bundled CJS modules that also declare `createRequire`. + banner: { + js: [ + `import { createRequire as __txeCreateRequire } from 'node:module';`, + `const require = __txeCreateRequire(import.meta.url);`, + ].join('\n'), + }, + // Strip whitespace and dead code, but do *not* rename identifiers — at least one of our bundled + // dependencies relies on global identifier names (we saw "position is not defined" when full + // `minify: true` was on). Whitespace+syntax minification still trims parse cost and is safe. + minifyWhitespace: true, + minifySyntax: true, + keepNames: true, + // External sourcemap keeps the runtime bundles small — V8 only parses the .js, and Node loads + // the .map lazily when a stack trace needs it. + sourcemap: 'external', + logLevel: 'info', + metafile: true, +}); + +// Dump the full metafile so we can audit chunk graph / imports separately from the build. +await writeFile('dest/metafile.json', JSON.stringify(result.metafile, null, 2)); + +const totalBytes = Object.values(result.metafile.outputs).reduce((sum, o) => sum + o.bytes, 0); +const ms = Date.now() - start; +// eslint-disable-next-line no-console +console.log(`Bundled TXE in ${ms}ms (${(totalBytes / 1024 / 1024).toFixed(1)} MiB total)`); + +// Surface the heaviest inputs per bundle. Pass `--inspect` on the command line to print. +if (process.argv.includes('--inspect')) { + for (const [outPath, out] of Object.entries(result.metafile.outputs)) { + if (!outPath.endsWith('.js')) { + continue; + } + // eslint-disable-next-line no-console + console.log(`\n${outPath} (${(out.bytes / 1024 / 1024).toFixed(1)} MiB) — top 40 contributors:`); + const inputs = Object.entries(out.inputs) + .sort(([, a], [, b]) => b.bytesInOutput - a.bytesInOutput) + .slice(0, 40); + for (const [path, info] of inputs) { + // eslint-disable-next-line no-console + console.log(` ${(info.bytesInOutput / 1024 / 1024).toFixed(2).padStart(6)} MiB ${path}`); + } + } +} + +enforceSizeLimits(result.metafile); diff --git a/yarn-project/txe/esbuild/plugins/external_native.mjs b/yarn-project/txe/esbuild/plugins/external_native.mjs new file mode 100644 index 000000000000..c3ca95facb25 --- /dev/null +++ b/yarn-project/txe/esbuild/plugins/external_native.mjs @@ -0,0 +1,11 @@ +/** + * Auto-externalizes any module whose resolved path is a .node native binding. Catches both + * direct `.node` imports and bare specifiers whose package main field points at a .node binary + * (e.g. `@napi-rs/snappy-linux-x64-gnu`). esbuild cannot bundle .node binaries. + */ +export const externalNativePlugin = { + name: 'external-native', + setup(build) { + build.onResolve({ filter: /\.node$/ }, args => ({ path: args.path, external: true })); + }, +}; diff --git a/yarn-project/txe/esbuild/plugins/protocol_contracts_eager_to_lazy.mjs b/yarn-project/txe/esbuild/plugins/protocol_contracts_eager_to_lazy.mjs new file mode 100644 index 000000000000..1684ec22b8e5 --- /dev/null +++ b/yarn-project/txe/esbuild/plugins/protocol_contracts_eager_to_lazy.mjs @@ -0,0 +1,42 @@ +/** + * Rewrites EAGER protocol-contract barrel imports (`class-registry`, `instance-registry`) to + * their `/lazy` siblings. The archiver, simulator, and aztec-node import these subpaths only + * for event classes (`ContractClassPublishedEvent`, `ContractInstancePublishedEvent`, …) which + * the `/lazy` barrel also re-exports — but the eager barrel statically pulls in the artifact + * JSONs at module init. Routing through `/lazy` keeps each artifact in its own per-contract + * dynamic chunk, loaded only when warmUp() runs. + * + * `fee-juice/index.js` is intentionally NOT rewritten: the eager barrel exposes synchronous + * utilities (`computeFeePayerBalanceStorageSlot`) the simulator calls on every public tx, and + * those need `storageLayout.balances.slot` from the loaded artifact. + * + * Two filters handle the two specifier forms seen by esbuild: + * - `@aztec/protocol-contracts/class-registry` — the package subpath, before esbuild does + * package.json `exports` resolution. We delegate back to esbuild via `build.resolve()` with + * `/lazy` appended so it picks the right per-package file. + * - `.../protocol-contracts/(dest|src)/(class|instance)-registry/index.(js|ts)` — the + * relative-path form after resolution, swapped to `lazy.(js|ts)` in place. + */ +export const protocolContractsEagerToLazyPlugin = { + name: 'protocol-contracts-event-subpath-stub', + setup(build) { + build.onResolve({ filter: /^@aztec\/protocol-contracts\/(class-registry|instance-registry)$/ }, async args => { + const result = await build.resolve(`${args.path}/lazy`, { + kind: args.kind, + importer: args.importer, + resolveDir: args.resolveDir, + pluginData: args.pluginData, + }); + return result.errors.length ? null : { path: result.path }; + }); + build.onResolve( + { filter: /protocol-contracts\/(dest|src)\/(class-registry|instance-registry)\/index\.(js|ts)$/ }, + args => { + if (args.path.endsWith('/lazy.js') || args.path.endsWith('/lazy.ts')) { + return null; + } + return { path: args.path.replace(/\/index\.(js|ts)$/, '/lazy.$1') }; + }, + ); + }, +}; diff --git a/yarn-project/txe/esbuild/plugins/redirects.mjs b/yarn-project/txe/esbuild/plugins/redirects.mjs new file mode 100644 index 000000000000..eb61a6af4f96 --- /dev/null +++ b/yarn-project/txe/esbuild/plugins/redirects.mjs @@ -0,0 +1,29 @@ +/** + * Generic resolver-redirect plugin. Each `rule` is `{filter, stub, importerContains?}`: + * - `filter`: regex matched against the specifier esbuild is trying to resolve. + * - `stub`: basename of a file under `esbuild/stubs/` to redirect to. + * - `importerContains` (optional): only fire if the importing file's path contains this + * substring. Used to gate relative-path filters (e.g. `./node_metrics.js`) so we don't + * accidentally hijack same-named files in another package. + * + * Rules are evaluated in order on each resolve callback; the first matching rule wins. + */ +export function redirectsPlugin(rules) { + const resolved = rules.map(rule => ({ + ...rule, + stubPath: new URL(`../stubs/${rule.stub}`, import.meta.url).pathname, + })); + return { + name: 'redirects', + setup(build) { + for (const { filter, importerContains, stubPath } of resolved) { + build.onResolve({ filter }, args => { + if (importerContains && !args.importer.includes(importerContains)) { + return null; + } + return { path: stubPath }; + }); + } + }, + }; +} diff --git a/yarn-project/txe/esbuild/plugins/size_guard.mjs b/yarn-project/txe/esbuild/plugins/size_guard.mjs new file mode 100644 index 000000000000..556801f4e1e2 --- /dev/null +++ b/yarn-project/txe/esbuild/plugins/size_guard.mjs @@ -0,0 +1,59 @@ +/** + * Post-build size guard. Catches unintended bundle growth without involving CI separately. + * + * Each entry pairs a regex against the output path with a `maxKB` cap and a `description` that + * shows up in the failure message. The build fails (exit 1) if any matching file exceeds its + * cap, or if the total bundle size exceeds `totalLimitMiB`. + * + * When a legitimate change pushes a chunk over its limit, raise the number AND append a one-line + * entry to the bump log so the history of size bumps stays auditable. + */ + +// Bump log: +// - 2026-05-27: initial limits. +export const sizeLimits = [ + // Shared chunks emitted by code-splitting; carry the simulator + PXE + world-state graph. + // Spikes here usually mean a heavy dep crept into the eager import path. + { pattern: /^dest\/chunk-.*\.js$/, maxKB: 2200, description: 'split chunk' }, + // Per-protocol-contract artifact chunks (loaded lazily via LazyProtocolContractsProvider). + { pattern: /^dest\/[A-Z][A-Za-z]+-[A-Z0-9]+\.js$/, maxKB: 800, description: 'contract artifact chunk' }, + // Tiny entry stubs that just re-export from the shared chunks. + { pattern: /^dest\/(worker|server)\.bundle\.js$/, maxKB: 8, description: 'entrypoint stub' }, + { pattern: /^dest\/bin\/index\.js$/, maxKB: 8, description: 'CLI entrypoint stub' }, +]; + +export const totalLimitMiB = 14; + +/** + * Validates a built esbuild `metafile` against the configured limits. Logs all violations then + * calls `process.exit(1)` if any were found. + */ +export function enforceSizeLimits(metafile) { + const violations = []; + let totalBytes = 0; + for (const [outPath, out] of Object.entries(metafile.outputs)) { + totalBytes += out.bytes; + for (const limit of sizeLimits) { + if (limit.pattern.test(outPath)) { + const sizeKB = out.bytes / 1024; + if (sizeKB > limit.maxKB) { + violations.push(` ${outPath}: ${sizeKB.toFixed(1)} KB > ${limit.maxKB} KB (${limit.description})`); + } + } + } + } + const totalMiB = totalBytes / 1024 / 1024; + if (totalMiB > totalLimitMiB) { + violations.push(` total: ${totalMiB.toFixed(2)} MiB > ${totalLimitMiB} MiB`); + } + if (violations.length === 0) { + return; + } + // eslint-disable-next-line no-console + console.error('\nBundle size guard tripped:\n' + violations.join('\n')); + // eslint-disable-next-line no-console + console.error( + '\nIf the new size is intentional, raise the corresponding limit in esbuild/size_guard.mjs and add a bump-log line.', + ); + process.exit(1); +} diff --git a/yarn-project/txe/esbuild/plugins/strip_artifact_debug.mjs b/yarn-project/txe/esbuild/plugins/strip_artifact_debug.mjs new file mode 100644 index 000000000000..03d9822cd976 --- /dev/null +++ b/yarn-project/txe/esbuild/plugins/strip_artifact_debug.mjs @@ -0,0 +1,30 @@ +import { readFile } from 'node:fs/promises'; + +/** + * Strips `file_map` and per-function `debug_symbols` from bundled contract artifact JSON files. + * Roughly 75% of each compiled artifact is sourcemaps / file maps used only by + * `getFunctionDebugMetadata` for error-trace enrichment - stripping these only + * results in TXE losing the capacity for private-function failure traces against the + * bundled protocol contracts or the SchnorrAccount artifact. + * + * Affects only the artifacts inlined into the bundle (protocol-contracts + SchnorrAccount). + * User contracts loaded at runtime during tests keep their full metadata. + */ +export const stripArtifactDebugPlugin = { + name: 'strip-artifact-debug', + setup(build) { + build.onLoad({ filter: /(protocol-contracts\/artifacts|accounts\/artifacts).*\.json$/ }, async args => { + const raw = await readFile(args.path, 'utf-8'); + const json = JSON.parse(raw); + // `ContractArtifactSchema` (yarn-project/stdlib/src/abi/abi.ts) requires `fileMap` to be a + // record and `debugSymbols` to be a string, so we zero them out instead of deleting them. + json.file_map = {}; + if (Array.isArray(json.functions)) { + for (const fn of json.functions) { + fn.debug_symbols = ''; + } + } + return { contents: JSON.stringify(json), loader: 'json' }; + }); + }, +}; diff --git a/yarn-project/txe/esbuild/stubs/archiver_stub.ts b/yarn-project/txe/esbuild/stubs/archiver_stub.ts new file mode 100644 index 000000000000..163b03df5cac --- /dev/null +++ b/yarn-project/txe/esbuild/stubs/archiver_stub.ts @@ -0,0 +1,17 @@ +import { throwStub } from './stub_helpers.js'; + +/* eslint-disable no-restricted-imports, import-x/no-relative-packages */ +export { ArchiverDataSourceBase } from '../../../archiver/dest/modules/data_source_base.js'; +export { ArchiverDataStoreUpdater } from '../../../archiver/dest/modules/data_store_updater.js'; +export { createArchiverDataStores } from '../../../archiver/dest/store/data_stores.js'; + +export function createArchiver(..._args: unknown[]): never { + throwStub('createArchiver'); +} + +export class L1ToL2MessagesNotReadyError extends Error { + constructor(message?: string) { + super(message); + this.name = 'L1ToL2MessagesNotReadyError'; + } +} diff --git a/yarn-project/txe/esbuild/stubs/aztec_node_node_metrics_stub.ts b/yarn-project/txe/esbuild/stubs/aztec_node_node_metrics_stub.ts new file mode 100644 index 000000000000..c43cbdea56ca --- /dev/null +++ b/yarn-project/txe/esbuild/stubs/aztec_node_node_metrics_stub.ts @@ -0,0 +1,20 @@ +import { throwStub } from './stub_helpers.js'; + +export class NodeMetrics { + // The constructor runs whenever an AztecNodeService is built, so it must stay a no-op. The + // recording methods below are only reached via sendTx / startSnapshotUpload, neither of which the + // TXE node exercises, so they throw if ever called. + constructor(_client: unknown, _name?: string) {} + + receivedTx(_durationMs: number, _isAccepted: boolean): void { + throwStub('NodeMetrics.receivedTx'); + } + + recordSnapshot(_durationMs: number): void { + throwStub('NodeMetrics.recordSnapshot'); + } + + recordSnapshotError(): void { + throwStub('NodeMetrics.recordSnapshotError'); + } +} diff --git a/yarn-project/txe/esbuild/stubs/aztec_node_sentinel_factory_stub.ts b/yarn-project/txe/esbuild/stubs/aztec_node_sentinel_factory_stub.ts new file mode 100644 index 000000000000..564e8a7fefe9 --- /dev/null +++ b/yarn-project/txe/esbuild/stubs/aztec_node_sentinel_factory_stub.ts @@ -0,0 +1,5 @@ +import { throwStub } from './stub_helpers.js'; + +export function createSentinel(..._args: unknown[]): never { + throwStub('createSentinel'); +} diff --git a/yarn-project/txe/esbuild/stubs/aztec_node_sentinel_stub.ts b/yarn-project/txe/esbuild/stubs/aztec_node_sentinel_stub.ts new file mode 100644 index 000000000000..35b62d412303 --- /dev/null +++ b/yarn-project/txe/esbuild/stubs/aztec_node_sentinel_stub.ts @@ -0,0 +1,7 @@ +import { throwStub } from './stub_helpers.js'; + +export class Sentinel { + constructor(..._args: unknown[]) { + throwStub('Sentinel'); + } +} diff --git a/yarn-project/txe/esbuild/stubs/bb_prover_stub.ts b/yarn-project/txe/esbuild/stubs/bb_prover_stub.ts new file mode 100644 index 000000000000..5b880aead156 --- /dev/null +++ b/yarn-project/txe/esbuild/stubs/bb_prover_stub.ts @@ -0,0 +1,19 @@ +import { throwStub } from './stub_helpers.js'; + +export class BBCircuitVerifier { + constructor(..._args: unknown[]) { + throwStub('BBCircuitVerifier'); + } +} + +export class BatchChonkVerifier { + constructor(..._args: unknown[]) { + throwStub('BatchChonkVerifier'); + } +} + +export class QueuedIVCVerifier { + constructor(..._args: unknown[]) { + throwStub('QueuedIVCVerifier'); + } +} diff --git a/yarn-project/txe/esbuild/stubs/bb_prover_test_stub.ts b/yarn-project/txe/esbuild/stubs/bb_prover_test_stub.ts new file mode 100644 index 000000000000..c284ff0b52ea --- /dev/null +++ b/yarn-project/txe/esbuild/stubs/bb_prover_test_stub.ts @@ -0,0 +1,2 @@ +/* eslint-disable no-restricted-imports, import-x/no-relative-packages */ +export { TestCircuitVerifier } from '../../../bb-prover/dest/test/test_verifier.js'; diff --git a/yarn-project/txe/esbuild/stubs/empty_stub.cjs b/yarn-project/txe/esbuild/stubs/empty_stub.cjs new file mode 100644 index 000000000000..0c066f42ff31 --- /dev/null +++ b/yarn-project/txe/esbuild/stubs/empty_stub.cjs @@ -0,0 +1,17 @@ +'use strict'; + +const handler = { + get(target, prop) { + if (prop === '__esModule') return true; + if (prop === 'default') return target; + if (typeof prop === 'symbol') return target[prop]; + if (!(prop in target)) { + target[prop] = new Proxy(function stubbed() { + throw new Error(`${String(prop)} is stubbed in this bundle; this code path should never run`); + }, handler); + } + return target[prop]; + }, +}; + +module.exports = new Proxy({}, handler); diff --git a/yarn-project/txe/esbuild/stubs/noble_bls12_stub.ts b/yarn-project/txe/esbuild/stubs/noble_bls12_stub.ts new file mode 100644 index 000000000000..212e72bae1ff --- /dev/null +++ b/yarn-project/txe/esbuild/stubs/noble_bls12_stub.ts @@ -0,0 +1,47 @@ +import { throwStub } from './stub_helpers.js'; + +const BLS12_FR_ORDER = 0x73eda753299d7d483339d80809a1d80553bda402fffe5bfeffffffff00000001n; +const BLS12_FP_ORDER = + 0x1a0111ea397fe69a4b1ba7b6434bacd764774b84f38512bf6730d2a0f6b0f6241eabfffeb153ffffb9feffffffffaaabn; + +function makeFieldStub(byteSize: number, order: bigint) { + return { + BYTES: byteSize, + ORDER: order, + MASK: (1n << BigInt(byteSize * 8)) - 1n, + ZERO: 0n, + ONE: 1n, + add: () => throwStub('bls12_381.add'), + sub: () => throwStub('bls12_381.sub'), + mul: () => throwStub('bls12_381.mul'), + div: () => throwStub('bls12_381.div'), + neg: () => throwStub('bls12_381.neg'), + sqr: () => throwStub('bls12_381.sqr'), + sqrt: () => throwStub('bls12_381.sqrt'), + pow: () => throwStub('bls12_381.pow'), + inv: () => throwStub('bls12_381.inv'), + eql: () => throwStub('bls12_381.eql'), + isValid: () => throwStub('bls12_381.isValid'), + is0: () => throwStub('bls12_381.is0'), + create: () => throwStub('bls12_381.create'), + fromBytes: () => throwStub('bls12_381.fromBytes'), + toBytes: () => throwStub('bls12_381.toBytes'), + }; +} + +const projectivePointStub = { + ZERO: { x: 0n, y: 0n, z: 0n, equals: (_other: unknown) => throwStub('bls12_381.equals') }, +}; + +// eslint-disable-next-line camelcase +export const bls12_381 = { + fields: { + Fr: makeFieldStub(32, BLS12_FR_ORDER), + Fp: makeFieldStub(48, BLS12_FP_ORDER), + }, + G1: { + CURVE: { Gx: 0n, Gy: 0n }, + ProjectivePoint: projectivePointStub, + weierstrassEquation: () => throwStub('bls12_381.weierstrassEquation'), + }, +}; diff --git a/yarn-project/txe/esbuild/stubs/noble_secp256k1_stub.ts b/yarn-project/txe/esbuild/stubs/noble_secp256k1_stub.ts new file mode 100644 index 000000000000..4429d94ceb75 --- /dev/null +++ b/yarn-project/txe/esbuild/stubs/noble_secp256k1_stub.ts @@ -0,0 +1,49 @@ +import { throwStub } from './stub_helpers.js'; + +const SECP_ORDER = 0xfffffffffffffffffffffffffffffffebaaedce6af48a03bbfd25e8cd0364141n; +const SECP_FIELD_ORDER = 0xfffffffffffffffffffffffffffffffffffffffffffffffffffffffefffffc2fn; + +const fieldStub = (byteSize: number, order: bigint) => ({ + BYTES: byteSize, + ORDER: order, + MASK: (1n << BigInt(byteSize * 8)) - 1n, + ZERO: 0n, + ONE: 1n, + add: () => throwStub('secp256k1.add'), + sub: () => throwStub('secp256k1.sub'), + mul: () => throwStub('secp256k1.mul'), + div: () => throwStub('secp256k1.div'), + neg: () => throwStub('secp256k1.neg'), + sqr: () => throwStub('secp256k1.sqr'), + sqrt: () => throwStub('secp256k1.sqrt'), + pow: () => throwStub('secp256k1.pow'), + inv: () => throwStub('secp256k1.inv'), + eql: () => throwStub('secp256k1.eql'), + isValid: () => throwStub('secp256k1.isValid'), + is0: () => throwStub('secp256k1.is0'), + create: () => throwStub('secp256k1.create'), + fromBytes: () => throwStub('secp256k1.fromBytes'), + toBytes: () => throwStub('secp256k1.toBytes'), +}); + +export const secp256k1 = { + CURVE: { n: SECP_ORDER, p: SECP_FIELD_ORDER, Gx: 0n, Gy: 0n }, + fields: { Fp: fieldStub(32, SECP_FIELD_ORDER), Fr: fieldStub(32, SECP_ORDER) }, + ProjectivePoint: { + ZERO: { x: 0n, y: 0n, z: 0n, equals: () => throwStub('secp256k1.equals') }, + BASE: { x: 0n, y: 0n, z: 0n }, + fromHex: () => throwStub('secp256k1.fromHex'), + fromAffine: () => throwStub('secp256k1.fromAffine'), + fromPrivateKey: () => throwStub('secp256k1.fromPrivateKey'), + }, + Signature: { fromCompact: () => throwStub('secp256k1.Signature.fromCompact') }, + utils: { + randomPrivateKey: () => throwStub('secp256k1.randomPrivateKey'), + isValidPrivateKey: () => throwStub('secp256k1.isValidPrivateKey'), + normPrivateKeyToScalar: () => throwStub('secp256k1.normPrivateKeyToScalar'), + }, + getPublicKey: () => throwStub('secp256k1.getPublicKey'), + sign: () => throwStub('secp256k1.sign'), + verify: () => throwStub('secp256k1.verify'), + getSharedSecret: () => throwStub('secp256k1.getSharedSecret'), +}; diff --git a/yarn-project/txe/esbuild/stubs/protocol_contract_eager_stub.ts b/yarn-project/txe/esbuild/stubs/protocol_contract_eager_stub.ts new file mode 100644 index 000000000000..cb0ff5c3b541 --- /dev/null +++ b/yarn-project/txe/esbuild/stubs/protocol_contract_eager_stub.ts @@ -0,0 +1 @@ +export {}; diff --git a/yarn-project/txe/esbuild/stubs/protocol_contracts_bundle_stub.ts b/yarn-project/txe/esbuild/stubs/protocol_contracts_bundle_stub.ts new file mode 100644 index 000000000000..c3c472f08599 --- /dev/null +++ b/yarn-project/txe/esbuild/stubs/protocol_contracts_bundle_stub.ts @@ -0,0 +1 @@ +export { LazyProtocolContractsProvider as BundledProtocolContractsProvider } from '@aztec/protocol-contracts/providers/lazy'; diff --git a/yarn-project/txe/esbuild/stubs/stub_helpers.ts b/yarn-project/txe/esbuild/stubs/stub_helpers.ts new file mode 100644 index 000000000000..561c4487613b --- /dev/null +++ b/yarn-project/txe/esbuild/stubs/stub_helpers.ts @@ -0,0 +1,20 @@ +/** + * Error thrown when a stubbed code path that must never run is invoked. Stubs exist only to satisfy + * an import shape for the bundled TXE; reaching one means the stubbed module is being used outside + * its intended context. + */ +export class StubInvocationError extends Error { + public override readonly name = 'StubInvocationError'; + + constructor(symbol: string) { + super(`${symbol} is stubbed in this bundle; this code path should never run`); + } +} + +/** Throws a {@link StubInvocationError} for `symbol`. Use for stubbed APIs that must never run. */ +export function throwStub(symbol: string): never { + throw new StubInvocationError(symbol); +} + +/** No-op for stubbed APIs that are exercised in the bundle but should intentionally do nothing. */ +export function noop(): void {} diff --git a/yarn-project/txe/esbuild/stubs/telemetry_start_stub.ts b/yarn-project/txe/esbuild/stubs/telemetry_start_stub.ts new file mode 100644 index 000000000000..c480184acb43 --- /dev/null +++ b/yarn-project/txe/esbuild/stubs/telemetry_start_stub.ts @@ -0,0 +1,15 @@ +/* eslint-disable no-restricted-imports, import-x/no-relative-packages */ +import { NoopTelemetryClient } from '../../../telemetry-client/dest/noop.js'; +import type { TelemetryClient } from '../../../telemetry-client/dest/telemetry.js'; + +export * from '../../../telemetry-client/dest/config.js'; + +const noopTelemetry: TelemetryClient = new NoopTelemetryClient(); + +export function getTelemetryClient(): TelemetryClient { + return noopTelemetry; +} + +export function initTelemetryClient(): Promise { + return Promise.resolve(noopTelemetry); +} diff --git a/yarn-project/txe/esbuild/stubs/telemetry_stub.ts b/yarn-project/txe/esbuild/stubs/telemetry_stub.ts new file mode 100644 index 000000000000..3ccd3129ae10 --- /dev/null +++ b/yarn-project/txe/esbuild/stubs/telemetry_stub.ts @@ -0,0 +1,39 @@ +/* eslint-disable no-restricted-imports, import-x/no-relative-packages */ +// eslint-disable-next-line import-x/no-extraneous-dependencies +import { ValueType } from '@opentelemetry/api'; + +import { noop } from './stub_helpers.js'; + +export * from '../../../telemetry-client/dest/telemetry.js'; +export * from '../../../telemetry-client/dest/noop.js'; +export * from '../../../telemetry-client/dest/with_tracer.js'; +export * from '../../../telemetry-client/dest/start.js'; +export * from '../../../telemetry-client/dest/otel_propagation.js'; +export * from '../../../telemetry-client/dest/prom_otel_adapter.js'; +export * from '../../../telemetry-client/dest/wrappers/fetch.js'; +export * from '../../../telemetry-client/dest/wrappers/l2_block_stream.js'; + +type MetricDefinition = { name: string; description: string; valueType: ValueType }; + +export class LmdbMetrics { + constructor(..._args: unknown[]) {} + recordDBMetrics = noop; + start = noop; + stop = noop; +} + +function makeMetricDefinition(prop: string): MetricDefinition { + return { + name: `aztec.stub.${prop.toLowerCase()}`, + description: 'TXE stub metric', + valueType: ValueType.INT, + }; +} + +export const Metrics: Record = new Proxy(Object.create(null), { + get: (_target, prop) => (typeof prop === 'string' ? makeMetricDefinition(prop) : undefined), +}); + +export const Attributes: Record = new Proxy(Object.create(null), { + get: (_target, prop) => (typeof prop === 'string' ? `aztec.stub.${prop.toLowerCase()}` : undefined), +}); diff --git a/yarn-project/txe/esbuild/stubs/world_state_stub.ts b/yarn-project/txe/esbuild/stubs/world_state_stub.ts new file mode 100644 index 000000000000..d5bdab6381f8 --- /dev/null +++ b/yarn-project/txe/esbuild/stubs/world_state_stub.ts @@ -0,0 +1,9 @@ +import { throwStub } from './stub_helpers.js'; + +export function createWorldState(..._args: unknown[]): never { + throwStub('createWorldState'); +} + +export function createWorldStateSynchronizer(..._args: unknown[]): never { + throwStub('createWorldStateSynchronizer'); +} diff --git a/yarn-project/txe/esbuild/stubs/zod_locales_stub.ts b/yarn-project/txe/esbuild/stubs/zod_locales_stub.ts new file mode 100644 index 000000000000..55f0a2e8fa9b --- /dev/null +++ b/yarn-project/txe/esbuild/stubs/zod_locales_stub.ts @@ -0,0 +1 @@ +export { default as en } from 'zod/v4/locales/en.js'; diff --git a/yarn-project/txe/eslint.config.js b/yarn-project/txe/eslint.config.js index 0331d0552f62..490437aab8e5 100644 --- a/yarn-project/txe/eslint.config.js +++ b/yarn-project/txe/eslint.config.js @@ -1,3 +1,7 @@ import config from '@aztec/foundation/eslint'; -export default config; +import { globalIgnores } from 'eslint/config'; + +// .cjs stubs are not part of the TS build (they're string-replaced by esbuild) and not +// listed in tsconfig, so typescript-eslint's project service refuses to parse them. +export default [...config, globalIgnores(['**/*.cjs'])]; diff --git a/yarn-project/txe/package.json b/yarn-project/txe/package.json index 913283c72c0a..52d5c630df3e 100644 --- a/yarn-project/txe/package.json +++ b/yarn-project/txe/package.json @@ -2,7 +2,12 @@ "name": "@aztec/txe", "version": "0.0.0", "type": "module", - "exports": "./dest/index.js", + "exports": { + "./server": { + "types": "./dest/rpc_server.d.ts", + "import": "./dest/server.bundle.js" + } + }, "bin": "./dest/bin/index.js", "typedocOptions": { "entryPoints": [ @@ -12,7 +17,7 @@ "tsconfig": "./tsconfig.json" }, "scripts": { - "build": "yarn clean && ../scripts/tsc.sh", + "build": "yarn clean && ../scripts/tsc.sh && node ./esbuild.config.mjs", "build:dev": "../scripts/tsc.sh --watch", "clean": "rm -rf ./dest .tsbuildinfo", "test": "NODE_NO_WARNINGS=1 node --experimental-vm-modules ../node_modules/.bin/jest --passWithNoTests --maxWorkers=${JEST_MAX_WORKERS:-8}", @@ -21,7 +26,8 @@ "check_txe_oracle_version": "node ./dest/bin/check_txe_oracle_version.js" }, "inherits": [ - "../package.common.json" + "../package.common.json", + "./package.local.json" ], "jest": { "moduleNameMapper": { @@ -67,6 +73,7 @@ "@aztec/aztec-node": "workspace:^", "@aztec/aztec.js": "workspace:^", "@aztec/bb-prover": "workspace:^", + "@aztec/bb.js": "workspace:^", "@aztec/constants": "workspace:^", "@aztec/foundation": "workspace:^", "@aztec/key-store": "workspace:^", @@ -76,6 +83,7 @@ "@aztec/simulator": "workspace:^", "@aztec/stdlib": "workspace:^", "@aztec/world-state": "workspace:^", + "msgpackr": "^1.11.2", "zod": "^4" }, "devDependencies": { @@ -83,17 +91,23 @@ "@types/jest": "^30.0.0", "@types/node": "^22.15.17", "@typescript/native-preview": "7.0.0-dev.20260113.1", + "esbuild": "^0.27.0", "jest": "^30.0.0", "jest-mock-extended": "^4.0.0", "ts-node": "^10.9.1", "typescript": "^5.3.3" }, "files": [ + "dest/bin/index.js", + "dest/worker.bundle.js", + "dest/server.bundle.js", + "dest/chunk-*.js", + "dest/*-*.js", + "dest/rpc_server.d.ts", "dest", "src", "!*.test.*" ], - "types": "./dest/index.d.ts", "engines": { "node": ">=20.10" } diff --git a/yarn-project/txe/package.local.json b/yarn-project/txe/package.local.json new file mode 100644 index 000000000000..bed13f9fdbee --- /dev/null +++ b/yarn-project/txe/package.local.json @@ -0,0 +1,5 @@ +{ + "scripts": { + "build": "yarn clean && ../scripts/tsc.sh && node ./esbuild.config.mjs" + } +} diff --git a/yarn-project/txe/src/bin/index.ts b/yarn-project/txe/src/bin/index.ts index 7c0b081931d2..a4d79a3543ab 100644 --- a/yarn-project/txe/src/bin/index.ts +++ b/yarn-project/txe/src/bin/index.ts @@ -2,7 +2,16 @@ import { createLogger } from '@aztec/aztec.js/log'; import { startHttpRpcServer } from '@aztec/foundation/json-rpc/server'; -import { createTXERpcServer } from '../index.js'; +import { createTXERpcServer } from '../rpc_server.js'; + +// Cap the native world-state thread pool before any import touches @aztec/world-state. +// `MAX_WORLD_STATE_THREADS` reads HARDWARE_CONCURRENCY at module load and defaults to 16; the +// pool's worker threads inherit `process.env`, so setting the cap here covers every worker. +// +// CAVEAT: HARDWARE_CONCURRENCY is process-global — bb prove/verify, the LMDB reader pool, and +// the world-state native thread pool all read it. Safe at 2 today because TXE never proves or +// verifies and only performs light tree updates; raise it if that ever changes. +process.env.HARDWARE_CONCURRENCY ??= '2'; /** * Create and start a new TXE HTTP Server @@ -23,7 +32,7 @@ async function main() { const logger = createLogger('txe:rpc'); logger.info(`Setting up TXE...`); - const txeServer = createTXERpcServer(logger); + const txeServer = await createTXERpcServer(logger); const { port } = await startHttpRpcServer(txeServer, { host: '127.0.0.1', port: TXE_PORT, diff --git a/yarn-project/txe/src/dispatcher_pool.ts b/yarn-project/txe/src/dispatcher_pool.ts new file mode 100644 index 000000000000..ddaa9783cf04 --- /dev/null +++ b/yarn-project/txe/src/dispatcher_pool.ts @@ -0,0 +1,314 @@ +import { getSchnorrAccountContractArtifact } from '@aztec/accounts/schnorr/lazy'; +import { BackendType, Barretenberg, BarretenbergSync } from '@aztec/bb.js'; +import type { Logger } from '@aztec/foundation/log'; +import { openEphemeralStore } from '@aztec/kv-store/lmdb-v2'; +import { LazyProtocolContractsProvider } from '@aztec/protocol-contracts/providers/lazy'; +import { ContractStore } from '@aztec/pxe/client/lazy'; +import { getContractClassFromArtifact } from '@aztec/stdlib/contract'; + +import { existsSync } from 'node:fs'; +import { cpus } from 'node:os'; +import { fileURLToPath } from 'node:url'; +import { Worker } from 'node:worker_threads'; + +import { TXE_REQUIRED_PROTOCOL_CONTRACTS } from './index.js'; +import type { TXEForeignCallInput } from './index.js'; +import type { ForeignCallResult } from './utils/encoding.js'; + +void Barretenberg.initSingleton({ backend: BackendType.Wasm, skipSrsInit: true, threads: 1 }); +void BarretenbergSync.initSingleton({ backend: BackendType.Wasm }); + +/** + * Opens a fresh LMDB in a tmp dir and writes the protocol contracts in + * {@link TXE_REQUIRED_PROTOCOL_CONTRACTS} plus the SchnorrAccount artifact, returning the + * directory path and the SchnorrAccount class id (hex). The store handle is intentionally kept + * alive: closing it would trigger the ephemeral-store cleanup hook and remove the tmp + * directory, so any worker that has not yet cloned would find it missing. + */ +export async function buildSharedContractStore(): Promise<{ dataDir: string; schnorrClassId: string }> { + const kvStore = await openEphemeralStore('txe-shared-contracts', undefined, 2); + const dataDir = kvStore.dataDirectory; + const contractStore = new ContractStore(kvStore); + const provider = new LazyProtocolContractsProvider(); + const [protocolContracts, schnorrArtifact] = await Promise.all([ + Promise.all(TXE_REQUIRED_PROTOCOL_CONTRACTS.map(name => provider.getProtocolContractArtifact(name))), + getSchnorrAccountContractArtifact(), + ]); + const schnorrClass = await getContractClassFromArtifact(schnorrArtifact); + await Promise.all([ + ...protocolContracts.flatMap(({ instance, artifact, contractClass }) => [ + contractStore.addContractArtifact(artifact, contractClass), + contractStore.addContractInstance(instance), + ]), + contractStore.addContractArtifact(schnorrArtifact, schnorrClass), + ]); + return { dataDir, schnorrClassId: schnorrClass.id.toString() }; +} + +/** + * Resolves `worker.bundle.js` whether this code is running unbundled (next to dispatcher_pool.js + * inside `dest/`) or bundled into `dest/bin/index.js` (one directory deeper). `import.meta.url` + * refers to whichever module the calling code actually lives in; we try both relative locations + * and use whichever exists. + */ +function resolveWorkerBundlePath(): URL { + const candidates = [new URL('./worker.bundle.js', import.meta.url), new URL('../worker.bundle.js', import.meta.url)]; + return candidates.find(u => existsSync(fileURLToPath(u))) ?? candidates[0]; +} + +interface SerializedError { + message: string; + name?: string; + stack?: string; +} + +type WorkerMessage = + | { type: 'result'; requestId: number; ok: true; value: ForeignCallResult } + | { type: 'result'; requestId: number; ok: false; error: SerializedError } + | { + type: 'memstat'; + sessions: number; + rss: number; + heapTotal: number; + heapUsed: number; + external: number; + arrayBuffers: number; + }; + +interface WorkerSlot { + worker: Worker; + sessions: Set; + inFlightRequestIds: Set; +} + +interface PendingRequest { + resolve: (value: ForeignCallResult) => void; + reject: (err: Error) => void; + workerIdx: number; +} + +function deserializeError(payload: SerializedError | undefined): Error { + if (!payload) { + return new Error('Worker returned an error with no payload'); + } + const err = new Error(payload.message); + if (payload.name) { + err.name = payload.name; + } + if (payload.stack) { + err.stack = payload.stack; + } + return err; +} + +export interface TXEDispatcherPoolOptions { + /** Number of worker threads */ + workers?: number; +} + +/** + * Main-thread router that owns a pool of TXE worker threads. Each worker runs its own + * {@link TXEDispatcher} and handles foreign-call requests for the sessions assigned to it. + * + * Routing is sticky by `session_id`: new sessions go to a freshly spawned worker (up to + * `maxWorkers`) or, once the cap is reached, to the existing worker with the fewest sessions. + * Each session's state — TXESession, native world state, KV stores — stays single-threaded + * within its worker; different sessions run in parallel across workers. + * + * The main thread builds a shared protocol-contracts LMDB once and passes its path via + * `workerData`. Workers clone the data file on demand instead of re-registering the contracts. + */ +export class TXEDispatcherPool { + private readonly workers: WorkerSlot[] = []; + private readonly sessionToWorker = new Map(); + private readonly pending = new Map(); + private nextRequestId = 0; + private readonly maxWorkers: number; + private readonly readyPromise: Promise; + private contractStoreSourceDir?: string; + private schnorrClassId?: string; + private workerPath?: URL; + + constructor( + private readonly logger: Logger, + opts: TXEDispatcherPoolOptions = {}, + ) { + // TXE still doesn't scale well beyond 16 workers due to contention + this.maxWorkers = Math.max(1, opts.workers ?? Math.min(cpus().length, 16)); + this.readyPromise = this.init(); + } + + /** Resolves once the shared protocol-contracts store is built. Workers spawn lazily on demand. */ + public ready(): Promise { + return this.readyPromise; + } + + private async init(): Promise { + const t0 = Date.now(); + const { dataDir, schnorrClassId } = await buildSharedContractStore(); + this.contractStoreSourceDir = dataDir; + this.schnorrClassId = schnorrClassId; + this.workerPath = resolveWorkerBundlePath(); + this.logger.debug(`TXE dispatcher pool ready (lazy spawn, cap=${this.maxWorkers})`, { + contractStoreSourceDir: dataDir, + schnorrClassId, + ms: Date.now() - t0, + }); + } + + /** + * Spawns a fresh worker and returns its index in `this.workers`. Caller must have awaited + * `init()` so the shared contract store path is set. Messages posted before the worker has + * finished loading are queued by Node's worker_threads transport. + */ + private spawnWorker(): number { + const workerIdx = this.workers.length; + const w = new Worker(this.workerPath!, { + workerData: { contractStoreSourceDir: this.contractStoreSourceDir, schnorrClassId: this.schnorrClassId }, + }); + const slot: WorkerSlot = { + worker: w, + sessions: new Set(), + inFlightRequestIds: new Set(), + }; + this.workers.push(slot); + w.on('message', (msg: WorkerMessage) => this.handleMessage(workerIdx, msg)); + w.on('error', err => this.handleWorkerError(workerIdx, err)); + w.on('exit', code => { + if (code !== 0) { + this.handleWorkerError(workerIdx, new Error(`Worker ${workerIdx} exited with code ${code}`)); + } + }); + this.logger.debug(`Spawning TXE worker ${workerIdx} on demand`, { + cap: this.maxWorkers, + poolSize: this.workers.length, + }); + return workerIdx; + } + + /** Routes a session-dispose request to the worker that owns the session. Fire-and-forget. */ + disposeSession(sessionId: number): void { + const workerIdx = this.sessionToWorker.get(sessionId); + if (workerIdx === undefined) { + throw new Error(`disposeSession: no worker mapped for session ${sessionId}`); + } + this.sessionToWorker.delete(sessionId); + const slot = this.workers[workerIdx]; + if (!slot) { + throw new Error(`disposeSession: worker ${workerIdx} (session ${sessionId}) missing from pool`); + } + slot.sessions.delete(sessionId); + slot.worker.postMessage({ type: 'dispose-session', sessionId }); + } + + // eslint-disable-next-line camelcase + async resolve_foreign_call(callData: TXEForeignCallInput): Promise { + // Make sure the shared contract store + worker bundle path are in place before we spawn. + await this.readyPromise; + const sessionId = callData.session_id; + let workerIdx = this.sessionToWorker.get(sessionId); + if (workerIdx === undefined) { + workerIdx = this.pickOrSpawnWorker(); + this.sessionToWorker.set(sessionId, workerIdx); + this.workers[workerIdx].sessions.add(sessionId); + this.logger.debug(`Routing new session ${sessionId} to worker ${workerIdx}`); + } + + const slot = this.workers[workerIdx]; + const requestId = this.nextRequestId++; + return new Promise((resolve, reject) => { + this.pending.set(requestId, { resolve, reject, workerIdx: workerIdx! }); + slot.inFlightRequestIds.add(requestId); + slot.worker.postMessage({ type: 'foreign-call', requestId, callData }); + }); + } + + /** Terminates all workers. Used by tests; production TXE exits the process directly. */ + async shutdown(): Promise { + const slots = this.workers.splice(0, this.workers.length); + await Promise.all(slots.map(s => s.worker.terminate())); + for (const [requestId, pending] of this.pending) { + pending.reject(new Error('TXE dispatcher pool was shut down')); + this.pending.delete(requestId); + } + } + + /** + * Returns the index of the worker that should handle a new session. Spawns a fresh worker if + * the pool hasn't reached its cap; otherwise picks the existing worker with the fewest + * assigned sessions. Always returns a valid index into `this.workers`. + */ + private pickOrSpawnWorker(): number { + if (this.workers.length < this.maxWorkers) { + return this.spawnWorker(); + } + let minIdx = 0; + let minCount = this.workers[0].sessions.size; + for (let i = 1; i < this.workers.length; i++) { + const count = this.workers[i].sessions.size; + if (count < minCount) { + minIdx = i; + minCount = count; + } + } + return minIdx; + } + + private handleMessage(workerIdx: number, msg: WorkerMessage): void { + switch (msg.type) { + case 'result': + this.handleResult(workerIdx, msg); + return; + case 'memstat': + this.logger.debug(`worker ${workerIdx} memstat`, { + worker: workerIdx, + sessions: msg.sessions, + rssMiB: Math.round(msg.rss / 1024 / 1024), + heapTotalMiB: Math.round(msg.heapTotal / 1024 / 1024), + heapUsedMiB: Math.round(msg.heapUsed / 1024 / 1024), + externalMiB: Math.round(msg.external / 1024 / 1024), + arrayBuffersMiB: Math.round(msg.arrayBuffers / 1024 / 1024), + }); + return; + } + } + + private handleResult(workerIdx: number, msg: WorkerMessage & { type: 'result' }): void { + const pending = this.pending.get(msg.requestId); + if (!pending) { + throw new Error(`handleResult: request ${msg.requestId} (worker ${workerIdx}) not in pending map`); + } + this.pending.delete(msg.requestId); + const slot = this.workers[workerIdx]; + if (!slot) { + throw new Error(`handleResult: worker ${workerIdx} (request ${msg.requestId}) missing from pool`); + } + slot.inFlightRequestIds.delete(msg.requestId); + if (msg.ok) { + pending.resolve(msg.value); + } else { + pending.reject(deserializeError(msg.error)); + } + } + + private handleWorkerError(workerIdx: number, err: Error): void { + this.logger.error(`TXE worker ${workerIdx} crashed; sessions assigned to it will fail`, err); + const slot = this.workers[workerIdx]; + if (!slot) { + throw new Error(`handleWorkerError: worker ${workerIdx} missing from pool (orig err: ${err.message})`); + } + for (const requestId of slot.inFlightRequestIds) { + const pending = this.pending.get(requestId); + if (!pending) { + throw new Error(`handleWorkerError: in-flight request ${requestId} (worker ${workerIdx}) not in pending map`); + } + pending.reject(err); + this.pending.delete(requestId); + } + slot.inFlightRequestIds.clear(); + for (const sessionId of slot.sessions) { + this.sessionToWorker.delete(sessionId); + } + slot.sessions.clear(); + } +} diff --git a/yarn-project/txe/src/index.ts b/yarn-project/txe/src/index.ts index d745b11df297..166debf47aa8 100644 --- a/yarn-project/txe/src/index.ts +++ b/yarn-project/txe/src/index.ts @@ -1,4 +1,3 @@ -import { SchnorrAccountContractArtifact } from '@aztec/accounts/schnorr'; import { type NoirCompiledContract, loadContractArtifact } from '@aztec/aztec.js/abi'; import { AztecAddress } from '@aztec/aztec.js/addresses'; import { @@ -7,12 +6,10 @@ import { } from '@aztec/aztec.js/contracts'; import { Fr } from '@aztec/aztec.js/fields'; import { PublicKeys, deriveKeys } from '@aztec/aztec.js/keys'; -import { createSafeJsonRpcServer } from '@aztec/foundation/json-rpc/server'; import type { Logger } from '@aztec/foundation/log'; -import { openTmpStore } from '@aztec/kv-store/lmdb-v2'; -import { protocolContractNames } from '@aztec/protocol-contracts'; -import { BundledProtocolContractsProvider } from '@aztec/protocol-contracts/providers/bundle'; -import { ContractStore } from '@aztec/pxe/server'; +import { cloneEphemeralStoreFrom } from '@aztec/kv-store/lmdb-v2'; +import type { ProtocolContractName } from '@aztec/protocol-contracts'; +import { ContractStore } from '@aztec/pxe/client/lazy'; import { computeArtifactHash } from '@aztec/stdlib/contract'; import type { ContractArtifactWithHash } from '@aztec/stdlib/contract'; import type { ApiSchemaFor } from '@aztec/stdlib/schemas'; @@ -24,6 +21,9 @@ import { readFile, readdir } from 'fs/promises'; import { join, parse } from 'path'; import { z } from 'zod'; +// Side-effect import: registers the msgpackr Fr extension for the bundled `Fr` class. Must +// be loaded before any `sendMessage` call. See msgpackr_fr_extension.ts for the why. +import './msgpackr_fr_extension.js'; import { type TXEOracleFunctionName, TXESession } from './txe_session.js'; import { type ForeignCallArgs, @@ -36,27 +36,57 @@ import { fromArray, fromSingle, toSingle, -} from './util/encoding.js'; +} from './utils/encoding.js'; + +// Protocol contracts TXE registers in its contract store. Only AuthRegistry is needed for the +// current test suites; add a contract here if a lookup against a `0x000…00X` address fails. +export const TXE_REQUIRED_PROTOCOL_CONTRACTS: ProtocolContractName[] = ['AuthRegistry']; const sessions = new Map(); -/* - * TXE typically has to load the same contract artifacts over and over again for multiple tests, - * so we cache them here to avoid loading from disk repeatedly. - * - * The in-flight map coalesces concurrent requests for the same cache key so that - * computeArtifactHash (very expensive) is only run once even under parallelism. +/** + * Cache + in-flight map pair. Lookup hits the cache, then awaits an in-flight `compute()` if one + * exists, otherwise starts one and stores it. Guarantees `compute()` runs at most once per `key` + * across concurrent callers, which matters because `computeArtifactHash` is expensive. */ -const TXEArtifactsCache = new Map< +class AsyncCache { + private readonly cache = new Map(); + private readonly inFlight = new Map>(); + + getOrCompute(key: K, compute: () => Promise): Promise { + const cached = this.cache.get(key); + if (cached !== undefined) { + return Promise.resolve(cached); + } + let pending = this.inFlight.get(key); + if (!pending) { + pending = (async () => { + try { + const value = await compute(); + this.cache.set(key, value); + return value; + } finally { + this.inFlight.delete(key); + } + })(); + this.inFlight.set(key, pending); + } + return pending; + } +} + +// Full deploys (artifact + computed instance), keyed by the full deploy context (contract + +// constructor args + publicKeys + salt + deployer). Hits on repeated identical deploys. +const TXEDeploymentsCache = new AsyncCache< string, { artifact: ContractArtifactWithHash; instance: ContractInstanceWithAddress } >(); -const TXEArtifactsCacheInFlight = new Map< - string, - Promise<{ artifact: ContractArtifactWithHash; instance: ContractInstanceWithAddress }> ->(); -type TXEForeignCallInput = { +// Loaded + hashed contract artifact, keyed by compiled-bytecode hash. Hits across deploys of the +// same contract when constructor args / salt / deployer differ. +const TXEArtifactsCache = new AsyncCache(); + +export type TXEForeignCallInput = { session_id: number; function: TXEOracleFunctionName; root_path: string; @@ -64,7 +94,7 @@ type TXEForeignCallInput = { inputs: ForeignCallArgs; }; -const TXEForeignCallInputSchema = zodFor()( +export const TXEForeignCallInputSchema = zodFor()( z.object({ // Nargo generates session_id as a u64, which may exceed Number.MAX_SAFE_INTEGER. // Zod 4's `.int()` enforces the safe-integer bound, so we drop it here and only require @@ -80,12 +110,56 @@ const TXEForeignCallInputSchema = zodFor()( }), ); -class TXEDispatcher { +export interface TXEDispatcherOptions { + /** + * Path to an LMDB directory holding the required protocol contracts (see + * {@link TXE_REQUIRED_PROTOCOL_CONTRACTS}) and the SchnorrAccount artifact. When set, the + * dispatcher clones this directory into a fresh tmpdir on first use instead of registering + * the contracts itself. + */ + contractStoreSourceDir: string; + /** + * Class id (hex) of the SchnorrAccount artifact pre-registered in the shared LMDB. When set, + * `#processAddAccountInputs` looks the artifact up from the cloned store instead of + * recomputing it via `getSchnorrAccountContractArtifact()` + `computeArtifactHash()`. + */ + schnorrClassId: string; +} + +export class TXEDispatcher { private contractStore!: ContractStore; + private readonly contractStoreSourceDir: string; + private readonly schnorrClassId: Fr; + + constructor( + private logger: Logger, + opts: TXEDispatcherOptions, + ) { + this.contractStoreSourceDir = opts.contractStoreSourceDir; + this.schnorrClassId = Fr.fromString(opts.schnorrClassId); + } - constructor(private logger: Logger) {} + /** + * Clones the pre-populated LMDB at `contractStoreSourceDir` into a fresh per-instance tmpdir + * on first use, so this dispatcher has a writable store already containing the required + * protocol contracts + SchnorrAccount. Idempotent — subsequent calls are no-ops. + */ + private async warmUp(): Promise { + if (this.contractStore) { + return; + } + const t0 = Date.now(); + const kvStore = await cloneEphemeralStoreFrom( + join(this.contractStoreSourceDir, 'data.mdb'), + 'txe-contracts', + undefined, + 2, + ); + this.contractStore = new ContractStore(kvStore); + this.logger.debug('Cloned shared protocol-contracts store', { totalMs: Date.now() - t0 }); + } - private fastHashFile(path: string) { + private fastHashFile(path: string): Promise { return new Promise(resolve => { const fd = createReadStream(path); const hash = createHash('sha1'); @@ -93,7 +167,7 @@ class TXEDispatcher { fd.on('end', function () { hash.end(); - resolve(hash.read()); + resolve(hash.read() as string); }); fd.pipe(hash); @@ -143,44 +217,29 @@ class TXEDispatcher { .map(arg => arg.toString()) .join('-')}-${publicKeysHash}-${salt}-${deployer}-${fileHash}`; - let instance; - let artifact: ContractArtifactWithHash; - - if (TXEArtifactsCache.has(cacheKey)) { - this.logger.debug(`Using cached artifact for ${cacheKey}`); - ({ artifact, instance } = TXEArtifactsCache.get(cacheKey)!); - } else { - if (!TXEArtifactsCacheInFlight.has(cacheKey)) { - this.logger.debug(`Loading compiled artifact ${artifactPath}`); - const compute = async () => { - const artifactJSON = JSON.parse(await readFile(artifactPath, 'utf-8')) as NoirCompiledContract; - const artifactWithoutHash = loadContractArtifact(artifactJSON); - const computedArtifact: ContractArtifactWithHash = { - ...artifactWithoutHash, - // Artifact hash is *very* expensive to compute, so we do it here once - // and the TXE contract data provider can cache it - artifactHash: await computeArtifactHash(artifactWithoutHash), - }; - this.logger.debug( - `Deploy ${computedArtifact.name} with initializer ${initializer}(${decodedArgs}) and public keys hash ${publicKeysHash.toString()}`, - ); - const computedInstance = await getContractInstanceFromInstantiationParams(computedArtifact, { - constructorArgs: decodedArgs, - skipArgsDecoding: true, - salt, - publicKeys, - constructorArtifact: initializer ? initializer : undefined, - deployer, - }); - const result = { artifact: computedArtifact, instance: computedInstance }; - TXEArtifactsCache.set(cacheKey, result); - TXEArtifactsCacheInFlight.delete(cacheKey); - return result; - }; - TXEArtifactsCacheInFlight.set(cacheKey, compute()); - } - ({ artifact, instance } = await TXEArtifactsCacheInFlight.get(cacheKey)!); - } + const { artifact, instance } = await TXEDeploymentsCache.getOrCompute(cacheKey, async () => { + this.logger.debug(`Loading compiled artifact ${artifactPath}`); + // Inner cache: artifact load + hash depends only on the compiled bytecode (`fileHash`), so + // subsequent deploys of the same contract — regardless of constructor args / deployer / + // salt — reuse the same `ContractArtifactWithHash`. + const computedArtifact = await TXEArtifactsCache.getOrCompute(fileHash, async () => { + const artifactJSON = JSON.parse(await readFile(artifactPath, 'utf-8')) as NoirCompiledContract; + const artifactWithoutHash = loadContractArtifact(artifactJSON); + return { ...artifactWithoutHash, artifactHash: await computeArtifactHash(artifactWithoutHash) }; + }); + this.logger.debug( + `Deploy ${computedArtifact.name} with initializer ${initializer}(${decodedArgs}) and public keys hash ${publicKeysHash.toString()}`, + ); + const computedInstance = await getContractInstanceFromInstantiationParams(computedArtifact, { + constructorArgs: decodedArgs, + skipArgsDecoding: true, + salt, + publicKeys, + constructorArtifact: initializer ? initializer : undefined, + deployer, + }); + return { artifact: computedArtifact, instance: computedInstance }; + }); inputs.splice(0, 1, artifact, instance, toSingle(secret)); } @@ -190,40 +249,29 @@ class TXEDispatcher { const cacheKey = `SchnorrAccountContract-${secret}`; - let artifact: ContractArtifactWithHash; - let instance; - - if (TXEArtifactsCache.has(cacheKey)) { - this.logger.debug(`Using cached artifact for ${cacheKey}`); - ({ artifact, instance } = TXEArtifactsCache.get(cacheKey)!); - } else { - if (!TXEArtifactsCacheInFlight.has(cacheKey)) { - const compute = async () => { - const keys = await deriveKeys(secret); - const args = [keys.publicKeys.ivpkM.x, keys.publicKeys.ivpkM.y]; - const computedArtifact: ContractArtifactWithHash = { - ...SchnorrAccountContractArtifact, - // Artifact hash is *very* expensive to compute, so we do it here once - // and the TXE contract data provider can cache it - artifactHash: await computeArtifactHash(SchnorrAccountContractArtifact), - }; - const computedInstance = await getContractInstanceFromInstantiationParams(computedArtifact, { - constructorArgs: args, - skipArgsDecoding: true, - salt: Fr.ONE, - publicKeys: keys.publicKeys, - constructorArtifact: 'constructor', - deployer: AztecAddress.ZERO, - }); - const result = { artifact: computedArtifact, instance: computedInstance }; - TXEArtifactsCache.set(cacheKey, result); - TXEArtifactsCacheInFlight.delete(cacheKey); - return result; - }; - TXEArtifactsCacheInFlight.set(cacheKey, compute()); + const { artifact, instance } = await TXEDeploymentsCache.getOrCompute(cacheKey, async () => { + const [artifactFromStore, classWithPreimage] = await Promise.all([ + this.contractStore.getContractArtifact(this.schnorrClassId), + this.contractStore.getContractClassWithPreimage(this.schnorrClassId), + ]); + if (!artifactFromStore || !classWithPreimage) { + throw new Error( + `SchnorrAccount not found in shared contract store at class id ${this.schnorrClassId.toString()}`, + ); } - ({ artifact, instance } = await TXEArtifactsCacheInFlight.get(cacheKey)!); - } + const computedArtifact = { ...artifactFromStore, artifactHash: classWithPreimage.artifactHash }; + const keys = await deriveKeys(secret); + const args = [keys.publicKeys.ivpkM.x, keys.publicKeys.ivpkM.y]; + const computedInstance = await getContractInstanceFromInstantiationParams(computedArtifact, { + constructorArgs: args, + skipArgsDecoding: true, + salt: Fr.ONE, + publicKeys: keys.publicKeys, + constructorArtifact: 'constructor', + deployer: AztecAddress.ZERO, + }); + return { artifact: computedArtifact, instance: computedInstance }; + }); inputs.splice(0, 0, artifact, instance); } @@ -235,17 +283,7 @@ class TXEDispatcher { if (!sessions.has(sessionId)) { this.logger.debug(`Creating new session ${sessionId}`); - if (!this.contractStore) { - const kvStore = await openTmpStore('txe-contracts'); - this.contractStore = new ContractStore(kvStore); - const provider = new BundledProtocolContractsProvider(); - for (const name of protocolContractNames) { - const { instance, artifact } = await provider.getProtocolContractArtifact(name); - await this.contractStore.addContractArtifact(artifact); - await this.contractStore.addContractInstance(instance); - } - this.logger.debug('Registered protocol contracts in shared contract store'); - } + await this.warmUp(); sessions.set(sessionId, await TXESession.init(this.contractStore)); } @@ -262,20 +300,31 @@ class TXEDispatcher { return await sessions.get(sessionId)!.processFunction(functionName, inputs); } + + /** + * Releases a session and its resources (per-session LMDB + `NativeWorldStateService`). + * Called by the dispatcher pool when nargo closes its TCP connection for a test (see + * `rpc_server.ts`'s socket tracker). No-op if the session was never created — that happens + * when nargo opens a connection but errors before sending a request. + */ + async disposeSession(sessionId: number): Promise { + const session = sessions.get(sessionId); + if (!session) { + return; + } + sessions.delete(sessionId); + await session.dispose(); + } +} + +/** Diagnostic-only: number of sessions currently held by this worker. */ +export function activeSessionCount(): number { + return sessions.size; } -const TXEDispatcherApiSchema: ApiSchemaFor = { +export const TXEDispatcherApiSchema: ApiSchemaFor = { // eslint-disable-next-line camelcase resolve_foreign_call: z.function({ input: z.tuple([TXEForeignCallInputSchema]), output: ForeignCallResultSchema }), + // disposeSession is invoked over IPC from the worker, not via RPC; required by ApiSchemaFor. + disposeSession: z.function({ input: z.tuple([z.number().nonnegative()]), output: z.void() }), }; - -/** - * Creates an RPC server that forwards calls to the TXE. - * @param logger - Logger to output to - * @returns A TXE RPC server. - */ -export function createTXERpcServer(logger: Logger) { - return createSafeJsonRpcServer(new TXEDispatcher(logger), TXEDispatcherApiSchema, { - http200OnError: true, - }); -} diff --git a/yarn-project/txe/src/msgpackr_fr_extension.ts b/yarn-project/txe/src/msgpackr_fr_extension.ts new file mode 100644 index 000000000000..5ac3bf745039 --- /dev/null +++ b/yarn-project/txe/src/msgpackr_fr_extension.ts @@ -0,0 +1,23 @@ +import { Fr } from '@aztec/foundation/curves/bn254'; + +import { addExtension } from 'msgpackr'; + +// `@aztec/native/msgpack_channel` registers a msgpackr extension that packs `Fr` instances to +// their raw 32-byte buffer for the C++ world-state side to deserialize. Because `@aztec/native` +// is externalized in TXE's esbuild config (it loads a `.node` binary that can't be bundled), +// the `Fr` class the extension is keyed on belongs to the *external* `@aztec/foundation` +// module instance — a different class identity than the `Fr` bundled into the TXE worker / +// bin. msgpackr's extension match uses `instanceof`, so bundled `Fr` instances slip through +// and msgpackr falls back to `Fr.toJSON()` (returns `Fr.toString()` — a `0x...` hex string). +// The C++ side then sees a string where a 32-byte binary was expected and throws +// `msgpack::type_error` whose `what()` is literally `"std::bad_cast"`. +// +// Fix: register the same extension under the BUNDLED `Fr` class identity. msgpackr's extension +// table holds an entry per `Class`, so both Fr classes now resolve to the same `write` callback +// and produce the same wire format. Imported as a side-effect from `./index.js` so both the +// worker (via `./worker.ts`) and the main-thread RPC server (via `./rpc_server.ts`) get the +// registration before any `sendMessage` runs. +addExtension({ + Class: Fr, + write: (fr: Fr) => fr.toBuffer(), +}); diff --git a/yarn-project/txe/src/oracle/txe_oracle_top_level_context.ts b/yarn-project/txe/src/oracle/txe_oracle_top_level_context.ts index 8338c84defae..2d9dc4d68222 100644 --- a/yarn-project/txe/src/oracle/txe_oracle_top_level_context.ts +++ b/yarn-project/txe/src/oracle/txe_oracle_top_level_context.ts @@ -80,13 +80,13 @@ import { collectNested, } from '@aztec/stdlib/tx'; import type { UInt64 } from '@aztec/stdlib/types'; -import { ForkCheckpoint } from '@aztec/world-state'; +import { ForkCheckpoint } from '@aztec/world-state/native'; import { DEFAULT_ADDRESS } from '../constants.js'; import type { TXEStateMachine } from '../state_machine/index.js'; -import type { TXEAccountStore } from '../util/txe_account_store.js'; -import { TXEPublicContractDataSource } from '../util/txe_public_contract_data_source.js'; import { getSingleTxBlockRequestHash, insertTxEffectIntoWorldTrees, makeTXEBlock } from '../utils/block_creation.js'; +import type { TXEAccountStore } from '../utils/txe_account_store.js'; +import { TXEPublicContractDataSource } from '../utils/txe_public_contract_data_source.js'; import type { ITxeExecutionOracle } from './interfaces.js'; export class TXEOracleTopLevelContext implements IMiscOracle, ITxeExecutionOracle { diff --git a/yarn-project/txe/src/rpc_server.ts b/yarn-project/txe/src/rpc_server.ts new file mode 100644 index 000000000000..f76668f1f336 --- /dev/null +++ b/yarn-project/txe/src/rpc_server.ts @@ -0,0 +1,87 @@ +import { createSafeJsonRpcServer } from '@aztec/foundation/json-rpc/server'; +import type { Logger } from '@aztec/foundation/log'; + +import type { Socket } from 'node:net'; + +import { TXEDispatcherPool, buildSharedContractStore } from './dispatcher_pool.js'; +import { TXEDispatcher, TXEDispatcherApiSchema } from './index.js'; + +/** + * Symbol used to tag an incoming TCP socket with the `session_id` it has been associated with. + * Hidden under a Symbol so we don't risk colliding with anything koa or http core adds. + */ +const SESSION_SYMBOL = Symbol('txeSessionId'); + +type TaggedSocket = Socket & { [SESSION_SYMBOL]?: number }; + +/** + * Creates the TXE RPC server. With `TXE_WORKERS=1` oracle calls run on the main thread (no + * worker_threads, no IPC overhead). With any other value oracle calls are + * routed to a pool of worker threads sized to that value, sticky by `session_id`. + * + * Each incoming TCP socket is tagged with the `session_id` of the first oracle call it carries — + * nargo uses one HTTP client per test, so the socket-to-session mapping is 1:1. When the socket + * closes (end of test), the dispatcher disposes the session and frees its world state + LMDB. + * + * Lives in its own module so the worker bundle does not pull in the HTTP server stack. + */ +export async function createTXERpcServer(logger: Logger) { + const workerCount = Number(process.env.TXE_WORKERS); + let dispatcher: TXEDispatcher | TXEDispatcherPool; + if (workerCount === 1) { + const { dataDir, schnorrClassId } = await buildSharedContractStore(); + dispatcher = new TXEDispatcher(logger, { contractStoreSourceDir: dataDir, schnorrClassId }); + } else { + dispatcher = new TXEDispatcherPool(logger, { + workers: Number.isFinite(workerCount) && workerCount > 1 ? workerCount : undefined, + }); + } + const server = createSafeJsonRpcServer(dispatcher, TXEDispatcherApiSchema, { + http200OnError: true, + middlewares: [ + async (ctx, next) => { + // The body parser runs further down the chain, so `ctx.request.body` is populated only + // after `next()` resolves. + await next(); + const socket = ctx.req.socket as TaggedSocket; + if (socket[SESSION_SYMBOL] !== undefined) { + return; + } + const body = (ctx.request as { body?: unknown }).body; + const sessionId = + body && typeof body === 'object' && 'params' in body + ? extractSessionId((body as { params: unknown }).params) + : undefined; + if (sessionId === undefined) { + return; + } + socket[SESSION_SYMBOL] = sessionId; + logger.debug(`Tagged socket with session`, { + sessionId, + remoteAddress: socket.remoteAddress, + remotePort: socket.remotePort, + }); + socket.once('close', () => { + logger.debug(`Disposing session on socket close`, { sessionId }); + void dispatcher.disposeSession(sessionId); + }); + }, + ], + }); + return server; +} + +// Extracts `session_id` from a JSON-RPC `params` array. Always a `number` because session_id +// comes off `JSON.parse`, which never produces BigInt; values above MAX_SAFE_INTEGER lose +// precision but still work as a Map key. +function extractSessionId(params: unknown): number | undefined { + if (!Array.isArray(params) || params.length === 0) { + return undefined; + } + const first = params[0]; + if (first && typeof first === 'object' && 'session_id' in first) { + const sid = (first as { session_id: unknown }).session_id; + return typeof sid === 'number' ? sid : undefined; + } + return undefined; +} diff --git a/yarn-project/txe/src/rpc_translator.ts b/yarn-project/txe/src/rpc_translator.ts index de4c456f42ea..f26ecc8478da 100644 --- a/yarn-project/txe/src/rpc_translator.ts +++ b/yarn-project/txe/src/rpc_translator.ts @@ -38,7 +38,7 @@ import { toArray, toForeignCallResult, toSingle, -} from './util/encoding.js'; +} from './utils/encoding.js'; const MAX_EVENT_LEN = 10; // This is MAX_MESSAGE_CONTENT_LEN - PRIVATE_EVENT_MSG_PLAINTEXT_RESERVED_FIELDS_LEN const MAX_PRIVATE_EVENTS_PER_TXE_QUERY = 5; diff --git a/yarn-project/txe/src/state_machine/index.ts b/yarn-project/txe/src/state_machine/index.ts index 746f79f0b9dd..69e2fad23220 100644 --- a/yarn-project/txe/src/state_machine/index.ts +++ b/yarn-project/txe/src/state_machine/index.ts @@ -9,7 +9,6 @@ import { L2Block, type L2TipsProvider } from '@aztec/stdlib/block'; import { Checkpoint, L1PublishedData, PublishedCheckpoint } from '@aztec/stdlib/checkpoint'; import type { AztecNode } from '@aztec/stdlib/interfaces/client'; import { CheckpointHeader } from '@aztec/stdlib/rollup'; -import { getPackageVersion } from '@aztec/stdlib/update-checker'; import { TXEArchiver } from './archiver.js'; import { DummyP2P } from './dummy_p2p_client.js'; @@ -19,6 +18,9 @@ import { TXESynchronizer } from './synchronizer.js'; const VERSION = 1; const CHAIN_ID = 1; +// Hardcoded so the bundled TXE doesn't try to read stdlib's package.json at a relative path +// that becomes invalid once bundled +const PACKAGE_VERSION = 'txe'; export class TXEStateMachine { constructor( @@ -59,7 +61,7 @@ export class TXEStateMachine { new TXEGlobalVariablesBuilder(), new TXEFeeProvider(), new MockEpochCache(), - getPackageVersion(), + PACKAGE_VERSION, new TestCircuitVerifier(), new TestCircuitVerifier(), undefined, diff --git a/yarn-project/txe/src/state_machine/synchronizer.ts b/yarn-project/txe/src/state_machine/synchronizer.ts index 8e217f4785d1..f6598d2ce8fe 100644 --- a/yarn-project/txe/src/state_machine/synchronizer.ts +++ b/yarn-project/txe/src/state_machine/synchronizer.ts @@ -18,7 +18,7 @@ export class TXESynchronizer implements WorldStateSynchronizer { constructor(public nativeWorldStateService: NativeWorldStateService) {} static async create() { - const nativeWorldStateService = await NativeWorldStateService.tmp(); + const nativeWorldStateService = await NativeWorldStateService.ephemeral(); return new this(nativeWorldStateService); } diff --git a/yarn-project/txe/src/txe_session.test.ts b/yarn-project/txe/src/txe_session.test.ts index 3dd5676b9ea2..876b833e8a73 100644 --- a/yarn-project/txe/src/txe_session.test.ts +++ b/yarn-project/txe/src/txe_session.test.ts @@ -9,6 +9,7 @@ describe('TXESession.processFunction', () => { beforeAll(() => { session = new TXESession( {} as any, // logger + {} as any, // sessionStore {} as any, // stateMachine {} as any, // oracleHandler {} as any, // contractStore diff --git a/yarn-project/txe/src/txe_session.ts b/yarn-project/txe/src/txe_session.ts index 32f674eeab30..16e39a4fa219 100644 --- a/yarn-project/txe/src/txe_session.ts +++ b/yarn-project/txe/src/txe_session.ts @@ -2,7 +2,8 @@ import { BlockNumber } from '@aztec/foundation/branded-types'; import { Fr } from '@aztec/foundation/curves/bn254'; import { type Logger, createLogger } from '@aztec/foundation/log'; import { KeyStore } from '@aztec/key-store'; -import { openTmpStore } from '@aztec/kv-store/lmdb-v2'; +import type { AztecAsyncKVStore } from '@aztec/kv-store'; +import { openEphemeralStore } from '@aztec/kv-store/lmdb-v2'; import { AddressStore, AnchorBlockStore, @@ -54,10 +55,10 @@ import { RPCTranslator } from './rpc_translator.js'; import { TXEArchiver } from './state_machine/archiver.js'; import { TXEStateMachine } from './state_machine/index.js'; import { TXE_ORACLE_VERSION_MAJOR, TXE_ORACLE_VERSION_MINOR } from './txe_oracle_version.js'; -import type { ForeignCallArgs, ForeignCallResult } from './util/encoding.js'; -import { TXEAccountStore } from './util/txe_account_store.js'; import { getSingleTxBlockRequestHash, insertTxEffectIntoWorldTrees, makeTXEBlock } from './utils/block_creation.js'; +import type { ForeignCallArgs, ForeignCallResult } from './utils/encoding.js'; import { makeTxEffect } from './utils/tx_effect_creation.js'; +import { TXEAccountStore } from './utils/txe_account_store.js'; /** * A TXE Session can be in one of four states, which change as the test progresses and different oracles are called. @@ -196,8 +197,11 @@ export class TXESession implements TXESessionStateHandler { private lastCallInfo: LastCallState = emptyLastCallState(); private txeOracleVersion: { major: number; minor: number } | undefined; + private disposed = false; + constructor( private logger: Logger, + private sessionStore: AztecAsyncKVStore, private stateMachine: TXEStateMachine, private oracleHandler: | IUtilityExecutionOracle @@ -221,8 +225,33 @@ export class TXESession implements TXESessionStateHandler { private nextBlockTimestamp: bigint, ) {} + /** + * Closes the per-session `txe-session` LMDB and the `NativeWorldStateService` . + * Called via IPC when the dispatcher detects the end of a test. Idempotent. + */ + async dispose(): Promise { + if (this.disposed) { + return; + } + this.disposed = true; + try { + await this.stateMachine.synchronizer.nativeWorldStateService.close(); + } catch (err) { + this.logger.warn(`Error closing native world state during session dispose`, err); + } + try { + await this.sessionStore.close(); + } catch (err) { + this.logger.warn(`Error closing session LMDB during dispose`, err); + } + } + static async init(contractStore: ContractStore) { - const store = await openTmpStore('txe-session'); + // Size LMDB's reader slots to the libuv pool (capped to 2 in bin/index.ts via + // HARDWARE_CONCURRENCY): each native LMDB read needs a libuv worker thread to run, so any + // slot beyond the pool size would sit idle while still consuming a semaphore + reader-table + // entry per session. + const store = await openEphemeralStore('txe-session', undefined, 2); const addressStore = new AddressStore(store); const privateEventStore = new PrivateEventStore(store); @@ -234,7 +263,6 @@ export class TXESession implements TXESessionStateHandler { const keyStore = new KeyStore(store); const accountStore = new TXEAccountStore(store); - // Create job coordinator and register staged stores const jobCoordinator = new JobCoordinator(store); jobCoordinator.registerStores([ capsuleStore, @@ -277,6 +305,7 @@ export class TXESession implements TXESessionStateHandler { return new TXESession( logger, + store, stateMachine, topLevelOracleHandler, contractStore, diff --git a/yarn-project/txe/src/util/encoding.ts b/yarn-project/txe/src/utils/encoding.ts similarity index 95% rename from yarn-project/txe/src/util/encoding.ts rename to yarn-project/txe/src/utils/encoding.ts index ea72550b40c3..03d66f8a7a38 100644 --- a/yarn-project/txe/src/util/encoding.ts +++ b/yarn-project/txe/src/utils/encoding.ts @@ -147,9 +147,10 @@ export function arrayOfArraysToBoundedVecOfArrays( const numFieldsToPad = maxLen * nestedArrayLength - flattenedStorage.length; - const flattenedStorageWithPadding = flattenedStorage.concat(Array(numFieldsToPad).fill(new Fr(0))); + // Pad with hex-encoded zeros so the result is uniformly `ForeignCallSingle` (string); raw Fr + // instances would JSON-serialize as `{ asBigInt: "0" }` and nargo would reject the response. + const flattenedStorageWithPadding = flattenedStorage.concat(toArray(Array(numFieldsToPad).fill(new Fr(0)))); - // At last we get the actual length of the BoundedVec and return the values. const len = toSingle(new Fr(bVecStorage.length)); return [flattenedStorageWithPadding, len]; } diff --git a/yarn-project/txe/src/util/expected_failure_error.ts b/yarn-project/txe/src/utils/expected_failure_error.ts similarity index 100% rename from yarn-project/txe/src/util/expected_failure_error.ts rename to yarn-project/txe/src/utils/expected_failure_error.ts diff --git a/yarn-project/txe/src/util/txe_account_store.ts b/yarn-project/txe/src/utils/txe_account_store.ts similarity index 100% rename from yarn-project/txe/src/util/txe_account_store.ts rename to yarn-project/txe/src/utils/txe_account_store.ts diff --git a/yarn-project/txe/src/util/txe_public_contract_data_source.ts b/yarn-project/txe/src/utils/txe_public_contract_data_source.ts similarity index 100% rename from yarn-project/txe/src/util/txe_public_contract_data_source.ts rename to yarn-project/txe/src/utils/txe_public_contract_data_source.ts diff --git a/yarn-project/txe/src/worker.ts b/yarn-project/txe/src/worker.ts new file mode 100644 index 000000000000..386070aadd60 --- /dev/null +++ b/yarn-project/txe/src/worker.ts @@ -0,0 +1,98 @@ +import { BackendType, Barretenberg, BarretenbergSync } from '@aztec/bb.js'; +import { type Logger, createLogger } from '@aztec/foundation/log'; + +import { parentPort, workerData } from 'node:worker_threads'; + +// Importing `./index.js` registers the msgpackr Fr extension transitively (via +// ./msgpackr_fr_extension.js); this must happen before any `sendMessage` call. +import { TXEDispatcher, type TXEDispatcherOptions, type TXEForeignCallInput, activeSessionCount } from './index.js'; + +// Seed both bb.js singletons with the WASM backend before any crypto call. `initSingleton` +// binds the singleton to whichever backend the first call requests, so this pre-empts the +// implicit `Barretenberg.initSingleton()` inside `poseidon2Hash` from `@aztec/foundation/crypto`. +// TXE only needs hashing (no proving, no verification), so WASM is sufficient and +// `skipSrsInit: true` skips the CRS load. `threads: 1` keeps the WASM backend on a single +// thread — additional threads would each spawn a nested worker_thread, multiplying memory cost +// per pool worker. +void Barretenberg.initSingleton({ backend: BackendType.Wasm, skipSrsInit: true, threads: 1 }); +void BarretenbergSync.initSingleton({ backend: BackendType.Wasm }); + +if (!parentPort) { + throw new Error('worker.ts must be loaded as a worker_thread'); +} + +const port = parentPort; +const logger: Logger = createLogger('txe:worker'); + +// The pool builds a template LMDB containing the protocol contracts + the SchnorrAccount +// artifact on the main thread and passes its data dir via `workerData`. The dispatcher clones +// that LMDB into a per-worker store on first use, so this worker gets a writable copy already +// populated. +const dispatcherOpts: TXEDispatcherOptions = { + contractStoreSourceDir: workerData.contractStoreSourceDir, + schnorrClassId: workerData.schnorrClassId, +}; +const dispatcher = new TXEDispatcher(logger, dispatcherOpts); + +interface ForeignCallRequest { + type: 'foreign-call'; + requestId: number; + callData: TXEForeignCallInput; +} + +interface DisposeSessionMessage { + type: 'dispose-session'; + sessionId: number; +} + +type IncomingMessage = ForeignCallRequest | DisposeSessionMessage; + +interface SerializedError { + message: string; + name?: string; + stack?: string; +} + +function serializeError(err: unknown): SerializedError { + if (err instanceof Error) { + return { message: err.message, name: err.name, stack: err.stack }; + } + return { message: String(err) }; +} + +// Periodic memstat for diagnostic builds — set TXE_WORKER_MEMSTAT=1 to enable. Posts JS-heap +// breakdown + active-session count back to the dispatcher so we can attribute RSS growth +// between V8 heap (would show in heapUsed) and native (LMDB / world-state / WASM). +if (process.env.TXE_WORKER_MEMSTAT === '1') { + setInterval(() => { + const m = process.memoryUsage(); + port.postMessage({ + type: 'memstat', + sessions: activeSessionCount(), + rss: m.rss, + heapTotal: m.heapTotal, + heapUsed: m.heapUsed, + external: m.external, + arrayBuffers: m.arrayBuffers, + }); + }, 2000).unref(); +} + +port.on('message', (msg: IncomingMessage) => { + switch (msg.type) { + case 'foreign-call': + void (async () => { + try { + const value = await dispatcher.resolve_foreign_call(msg.callData); + port.postMessage({ type: 'result', requestId: msg.requestId, ok: true, value }); + } catch (err) { + port.postMessage({ type: 'result', requestId: msg.requestId, ok: false, error: serializeError(err) }); + } + })(); + return; + case 'dispose-session': + // Fire-and-forget; the main thread does not wait for confirmation. + void dispatcher.disposeSession(msg.sessionId).catch(err => logger.warn(`disposeSession failed`, err)); + return; + } +}); diff --git a/yarn-project/world-state/src/native/native_world_state.test.ts b/yarn-project/world-state/src/native/native_world_state.test.ts index 5c4437ca7386..5acafa7d67b6 100644 --- a/yarn-project/world-state/src/native/native_world_state.test.ts +++ b/yarn-project/world-state/src/native/native_world_state.test.ts @@ -961,9 +961,6 @@ describe('NativeWorldState', () => { await expect(delayedFork.getSiblingPath(MerkleTreeId.NULLIFIER_TREE, 0n)).rejects.toThrow('Fork not found'); await sleep(closeDelayMs * 3); - const closePromise = (delayedFork as any).closePromise; - expect(closePromise).toBeDefined(); - await closePromise; expect(warnSpy).not.toHaveBeenCalled(); expect((ws as any).instance.queues.has(forkId)).toBe(false); diff --git a/yarn-project/world-state/src/native/native_world_state.ts b/yarn-project/world-state/src/native/native_world_state.ts index fd799993f3b1..4c69360b551c 100644 --- a/yarn-project/world-state/src/native/native_world_state.ts +++ b/yarn-project/world-state/src/native/native_world_state.ts @@ -44,6 +44,29 @@ export const WORLD_STATE_DB_VERSION = 2; // The initial version export const WORLD_STATE_DIR = 'world_state'; +const DEFAULT_TMP_TREE_MAP_SIZE_KB = 10 * 1024 * 1024; + +/** + * Sets up a fresh `mkdtemp` directory + default `WorldStateTreeMapSizes` shared by both + * the `.tmp` (fsync-on) and `.ephemeral` (fsync-off) factories. Returns the raw tmpdir, + * the tree map sizes, and the package logger. + */ +async function createTmpWorldStateDir( + bindings?: LoggerBindings, +): Promise<{ dataDir: string; wsTreeMapSizes: WorldStateTreeMapSizes; log: Logger }> { + const log = createLogger('world-state:database', bindings); + const dataDir = await mkdtemp(join(tmpdir(), 'aztec-world-state-')); + const wsTreeMapSizes: WorldStateTreeMapSizes = { + archiveTreeMapSizeKb: DEFAULT_TMP_TREE_MAP_SIZE_KB, + nullifierTreeMapSizeKb: DEFAULT_TMP_TREE_MAP_SIZE_KB, + noteHashTreeMapSizeKb: DEFAULT_TMP_TREE_MAP_SIZE_KB, + messageTreeMapSizeKb: DEFAULT_TMP_TREE_MAP_SIZE_KB, + publicDataTreeMapSizeKb: DEFAULT_TMP_TREE_MAP_SIZE_KB, + }; + log.debug(`Created temporary world state database at: ${dataDir} (map size ${DEFAULT_TMP_TREE_MAP_SIZE_KB} KB)`); + return { dataDir, wsTreeMapSizes, log }; +} + export class NativeWorldStateService implements MerkleTreeDatabase { protected initialHeader: BlockHeader | undefined; // This is read heavily and only changes when data is persisted, so we cache it @@ -57,6 +80,11 @@ export class NativeWorldStateService implements MerkleTreeDatabase { private readonly cleanup = () => Promise.resolve(), ) {} + /** + * Opens a persistent world state at `dataDir`. Goes through `DatabaseVersionManager` so the + * caller's rollup address is bound to the on-disk schema and incompatible versions surface + * loudly. The LMDB envs commit with full fsync. + */ static async new( rollupAddress: EthAddress, dataDir: string, @@ -68,14 +96,22 @@ export class NativeWorldStateService implements MerkleTreeDatabase { ): Promise { const log = createLogger('world-state:database', bindings); const worldStateDirectory = join(dataDir, WORLD_STATE_DIR); - // Create a version manager to handle versioning const versionManager = new DatabaseVersionManager({ schemaVersion: WORLD_STATE_DB_VERSION, rollupAddress, dataDirectory: worldStateDirectory, - onOpen: (dir: string) => { - return Promise.resolve(new NativeWorldState(dir, wsTreeMapSizes, genesis, instrumentation, bindings)); - }, + onOpen: (dir: string) => + Promise.resolve( + new NativeWorldState( + dir, + wsTreeMapSizes, + genesis, + instrumentation, + bindings, + undefined, + /*ephemeral=*/ false, + ), + ), }); const [instance] = await versionManager.open(); @@ -90,6 +126,15 @@ export class NativeWorldStateService implements MerkleTreeDatabase { return worldState; } + /** + * Opens a world state in a fresh tmpdir with full fsync semantics. Use when you need the + * on-disk file to remain crash-recoverable (e.g. for snapshot/backup tests) but don't + * want a persistent dataDir. Pass `cleanupTmpDir=false` to keep the directory after + * close for inspection. + * + * If you don't care about crash-recoverability — i.e. you just want a fast scratch + * database for tests — use {@link ephemeral} instead. + */ static async tmp( rollupAddress = EthAddress.ZERO, cleanupTmpDir = true, @@ -97,19 +142,7 @@ export class NativeWorldStateService implements MerkleTreeDatabase { instrumentation = new WorldStateInstrumentation(getTelemetryClient()), bindings?: LoggerBindings, ): Promise { - const log = createLogger('world-state:database', bindings); - const dataDir = await mkdtemp(join(tmpdir(), 'aztec-world-state-')); - const dbMapSizeKb = 10 * 1024 * 1024; - const worldStateTreeMapSizes: WorldStateTreeMapSizes = { - archiveTreeMapSizeKb: dbMapSizeKb, - nullifierTreeMapSizeKb: dbMapSizeKb, - noteHashTreeMapSizeKb: dbMapSizeKb, - messageTreeMapSizeKb: dbMapSizeKb, - publicDataTreeMapSizeKb: dbMapSizeKb, - }; - log.debug(`Created temporary world state database at: ${dataDir} with tree map size: ${dbMapSizeKb}`); - - // pass a cleanup callback because process.on('beforeExit', cleanup) does not work under Jest + const { dataDir, wsTreeMapSizes, log } = await createTmpWorldStateDir(bindings); const cleanup = async () => { if (cleanupTmpDir) { await rm(dataDir, { recursive: true, force: true, maxRetries: 3 }); @@ -118,8 +151,45 @@ export class NativeWorldStateService implements MerkleTreeDatabase { log.debug(`Leaving temporary world state database: ${dataDir}`); } }; + return this.new(rollupAddress, dataDir, wsTreeMapSizes, genesis, instrumentation, bindings, cleanup); + } - return this.new(rollupAddress, dataDir, worldStateTreeMapSizes, genesis, instrumentation, bindings, cleanup); + /** + * Opens a fully-ephemeral world state. The directory is created in `os.tmpdir()`, the LMDB + * envs open with `MDB_NOSYNC | MDB_NOMETASYNC` so commits never block on fsync, and the + * directory is removed on dispose. A crash mid-write leaves the env unrecoverable. + * + * For unit tests and other isolated runs. Use {@link tmp} when you need fsync semantics in a + * tmp dir, and {@link new} for a persistent store. Skips {@link DatabaseVersionManager} — + * there is no on-disk schema to bind to and no rollup address is taken. + */ + static async ephemeral( + genesis: GenesisData = EMPTY_GENESIS_DATA, + instrumentation = new WorldStateInstrumentation(getTelemetryClient()), + bindings?: LoggerBindings, + ): Promise { + const { dataDir, wsTreeMapSizes, log } = await createTmpWorldStateDir(bindings); + const cleanup = async () => { + await rm(dataDir, { recursive: true, force: true, maxRetries: 3 }); + log.debug(`Deleted ephemeral world state database: ${dataDir}`); + }; + const instance = new NativeWorldState( + join(dataDir, WORLD_STATE_DIR), + wsTreeMapSizes, + genesis, + instrumentation, + bindings, + undefined, + /*ephemeral=*/ true, + ); + const worldState = new this(instance, instrumentation, log, genesis, cleanup); + try { + await worldState.init(); + } catch (e) { + log.error(`Error initializing ephemeral world state: ${e}`); + throw e; + } + return worldState; } protected async init() { diff --git a/yarn-project/world-state/src/native/native_world_state_instance.ts b/yarn-project/world-state/src/native/native_world_state_instance.ts index c4016ba1e477..2d9a7b4658bd 100644 --- a/yarn-project/world-state/src/native/native_world_state_instance.ts +++ b/yarn-project/world-state/src/native/native_world_state_instance.ts @@ -59,12 +59,13 @@ export class NativeWorldState implements NativeWorldStateInstance { private readonly instrumentation: WorldStateInstrumentation, bindings?: LoggerBindings, private readonly log: Logger = createLogger('world-state:database', bindings), + private readonly ephemeral: boolean = false, ) { const threads = Math.min(cpus().length, MAX_WORLD_STATE_THREADS); log.info( `Creating world state data store at directory ${dataDir} with map sizes ${JSON.stringify( wsTreeMapSizes, - )} and ${threads} threads.`, + )} and ${threads} threads (ephemeral=${ephemeral}).`, ); const prefilledPublicDataBufferArray = genesis.prefilledPublicData.map(d => [ d.slot.toBuffer(), @@ -94,6 +95,7 @@ export class NativeWorldState implements NativeWorldStateInstance { [MerkleTreeId.ARCHIVE]: wsTreeMapSizes.archiveTreeMapSizeKb, }, threads, + ephemeral, ); this.instance = new MsgpackChannel(ws); // Manually create the queue for the canonical fork @@ -112,6 +114,7 @@ export class NativeWorldState implements NativeWorldStateInstance { this.instrumentation, this.log.getBindings(), this.log, + this.ephemeral, ); } diff --git a/yarn-project/yarn.lock b/yarn-project/yarn.lock index 4c76b9a61b4f..7b463d0a45e8 100644 --- a/yarn-project/yarn.lock +++ b/yarn-project/yarn.lock @@ -2149,6 +2149,7 @@ __metadata: "@aztec/aztec-node": "workspace:^" "@aztec/aztec.js": "workspace:^" "@aztec/bb-prover": "workspace:^" + "@aztec/bb.js": "workspace:^" "@aztec/constants": "workspace:^" "@aztec/foundation": "workspace:^" "@aztec/key-store": "workspace:^" @@ -2162,8 +2163,10 @@ __metadata: "@types/jest": "npm:^30.0.0" "@types/node": "npm:^22.15.17" "@typescript/native-preview": "npm:7.0.0-dev.20260113.1" + esbuild: "npm:^0.27.0" jest: "npm:^30.0.0" jest-mock-extended: "npm:^4.0.0" + msgpackr: "npm:^1.11.2" ts-node: "npm:^10.9.1" typescript: "npm:^5.3.3" zod: "npm:^4"