diff --git a/CHANGELOG.md b/CHANGELOG.md index 11ac5cde..a8ce3945 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -10,6 +10,24 @@ order; the entry is the editorial text of the version's GitHub Release. The 1.7 maintenance line continues in `CHANGELOG.md` on the `release/1.7` branch. This project follows semantic versioning. +## [2.2.0] + +Moves the filler's SHIP reader onto a published package. + +### Upgrading + +- Image `ghcr.io/atomicassets/atomicassets-api:2.2.0`. The `2.2` and `latest` tags move to it. +- The migration set is unchanged from 2.0.0, so the filler performs no database work on boot. +- A heartbeat runs every 30 seconds and the socket is torn down after 300 seconds without a message or a pong. The previous client had neither and hung on a half-open connection. (#175) +- Reconnect backoff runs 5 seconds to 60 seconds instead of a fixed 5 second retry. (#175) +- A block, trace or delta payload the node serves empty escalates to a reconnect instead of pausing the queue permanently. Empty payloads at `block_num <= 1` only warn. (#175) +- A failure on the prepare path rejects into an unhandled rejection, which the filler turns into `process.exit(1)` and a supervisor restart. The previous client paused the queue until the stall watchdog fired. (#175) +- `ship_min_block_confirmation`, `ds_ship_threads`, `ship_prefetch_blocks`, `ship_ds_queue_size` and `ship_max_blocks_queue` keep their meaning. Heartbeat and idle timeout are not configurable in `readers.config.json`. (#175) + +### Other changes + +- The filler's SHIP reader is the `@atomichub/antelope-ship-utils` package (repository `atomicassets/antelope-ship-utils`) rather than an in-tree client. The in-tree reader, its deserializer worker and the direct `ws` and `node-worker-threads-pool` dependencies are gone, and the block-shape types are re-exported from the package. (#175) + ## [2.1.0] - 2026-08-10 Splits the Redis pub/sub subscriber onto its own optional endpoint. diff --git a/knip.json b/knip.json index d8f5425b..d3f54273 100644 --- a/knip.json +++ b/knip.json @@ -2,7 +2,6 @@ "$schema": "https://unpkg.com/knip@5/schema.json", "entry": [ "src/bin/**/*.ts", - "src/workers/**/*.ts", "src/scripts/**/*.ts", "src/**/*.test.ts" ], diff --git a/package.json b/package.json index 24552430..2e3b2525 100644 --- a/package.json +++ b/package.json @@ -52,6 +52,7 @@ "db:migrate:up": "pnpm build && node --trace-warnings ./build/bin/migrate.js" }, "dependencies": { + "@atomichub/antelope-ship-utils": "^1.0.0", "@atomichub/atomicassets": "^2.0.2", "@wharfkit/antelope": "^1.1.1", "async-exit-hook": "^2.0.1", @@ -64,15 +65,13 @@ "express-rate-limit": "^8.5.2", "iovalkey": "^0.3.3", "moize": "^6.1.1", - "node-worker-threads-pool": "^1.5.1", "p-queue": "^9.3.1", "pg": "^8.22.0", "prom-client": "^15.1.3", "rate-limit-redis": "^5.0.0", "socket.io": "^4.8.3", "swagger-ui-express": "^5.0.1", - "winston": "^3.7.2", - "ws": "^8.21.1" + "winston": "^3.7.2" }, "devDependencies": { "@eslint/js": "^10.0.1", @@ -93,7 +92,6 @@ "@types/sinon": "^21.0.0", "@types/supertest": "^7.2.1", "@types/swagger-ui-express": "^4.1.6", - "@types/ws": "^8.18.1", "c8": "^11.0.0", "chai": "^6.2.2", "chai-as-promised": "^8.0.2", diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index e9262960..5b5d7bdf 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -11,6 +11,9 @@ importers: .: dependencies: + '@atomichub/antelope-ship-utils': + specifier: ^1.0.0 + version: 1.0.0 '@atomichub/atomicassets': specifier: ^2.0.2 version: 2.0.2 @@ -47,9 +50,6 @@ importers: moize: specifier: ^6.1.1 version: 6.1.7 - node-worker-threads-pool: - specifier: ^1.5.1 - version: 1.5.1 p-queue: specifier: ^9.3.1 version: 9.3.1 @@ -71,9 +71,6 @@ importers: winston: specifier: ^3.7.2 version: 3.19.0 - ws: - specifier: ^8.21.1 - version: 8.21.1 devDependencies: '@eslint/js': specifier: ^10.0.1 @@ -129,9 +126,6 @@ importers: '@types/swagger-ui-express': specifier: ^4.1.6 version: 4.1.8 - '@types/ws': - specifier: ^8.18.1 - version: 8.18.1 c8: specifier: ^11.0.0 version: 11.0.0 @@ -171,6 +165,10 @@ importers: packages: + '@atomichub/antelope-ship-utils@1.0.0': + resolution: {integrity: sha512-8w0AEF4DAfx0ZdiENb4NI3mJdNGJI3WI0qvRvOXcm7tEG8YlxWFZUEQn3OrYqFGpO3BxIVrlBDifyS9Cr/jhSA==} + engines: {node: '>=22'} + '@atomichub/atomicassets@2.0.2': resolution: {integrity: sha512-ae6L8+SDhQ9dtDozu/A8Qwkax06V80BCVrqkT180FWhG9gDkIWvkYREN4nE7OsWeil6Lf42vVlDdfRl1gzjuag==} engines: {node: '>=20'} @@ -1351,6 +1349,9 @@ packages: resolution: {integrity: sha512-aIL5Fx7mawVa300al2BnEE4iNvo1qETxLrPI/o05L7z6go7fCw1J6EQmbK4FmJ2AS7kgVF/KEZWufBfdClMcPg==} engines: {node: '>= 0.6'} + eventemitter3@4.0.7: + resolution: {integrity: sha512-8guHBZCwKnFhYdHr2ysuRWErTwhoN2X8XELRlrRwpmfeY2jjuUN4taQMsULKUVo1K4DvZl+0pgfyoysHxvmvEw==} + eventemitter3@5.0.4: resolution: {integrity: sha512-mlsTRyGaPBjPedk6Bvw+aqbsXDtoAyAzm5MO7JgU+yVRyMQ5O8bD4Kcci7BS85f93veegeCPkL8R4GLClnjLFw==} @@ -1968,6 +1969,10 @@ packages: resolution: {integrity: sha512-Q6Bekk5wpzW5qIyUP4gdMEujObYstZl6DMMOSenwBvV0BlE5LkDwkjs5yHbZmdCEq2o4RJx4tE1vwxFVf2FG1w==} engines: {node: '>=16.17'} + p-finally@1.0.0: + resolution: {integrity: sha512-LICb2p9CB7FS+0eR1oqWnHhp0FljGLZCWBE9aix0Uye9W8LTQPwMTYVGWQWIw9RdQiDg4+epXQODwIYJtSJaow==} + engines: {node: '>=4'} + p-limit@3.1.0: resolution: {integrity: sha512-TYOanM3wGwNGsZN2cVTYPArw454xnXj5qmWF1bEoAc4+cU/ol7GVh7odevjp1FNHduHc3KZMcFduxU5Xc6uJRQ==} engines: {node: '>=10'} @@ -1976,10 +1981,18 @@ packages: resolution: {integrity: sha512-LaNjtRWUBY++zB5nE/NwcaoMylSPk+S+ZHNB1TzdbMJMny6dynpAGt7X/tl/QYq3TIeE6nxHppbo2LGymrG5Pw==} engines: {node: '>=10'} + p-queue@6.6.2: + resolution: {integrity: sha512-RwFpb72c/BhQLEXIZ5K2e+AhgNVmIejGlTgiB9MzZ0e93GRvqZ7uSi0dvRF7/XIXDeNkra2fNHBxTyPDGySpjQ==} + engines: {node: '>=8'} + p-queue@9.3.1: resolution: {integrity: sha512-POWdiIPmsUPGwb4FeQ4OBg46aqmcInSWe45CKDsGHiOBiVQM9chqfQTuqhuTzcg2Vz9faTI65at0KkVyVEiCHw==} engines: {node: '>=20'} + p-timeout@3.2.0: + resolution: {integrity: sha512-rhIwUycgwwKcP9yTOOFK/AKsAopjjCakVqLHePO3CC6Mir1Z99xT+R63jZxAT5lFZLa2inS5h+ZS2GvR99/FBg==} + engines: {node: '>=8'} + p-timeout@6.1.4: resolution: {integrity: sha512-MyIV3ZA/PmyBN/ud8vV9XzwTrNtR4jFrObymZYnZqMmW0zA8Z17vnT0rBgFE/TlohB+YCHqXMgZzb3Csp49vqg==} engines: {node: '>=14.16'} @@ -2602,6 +2615,17 @@ packages: snapshots: + '@atomichub/antelope-ship-utils@1.0.0': + dependencies: + '@wharfkit/antelope': 1.2.0 + common-tags: 1.8.2 + node-worker-threads-pool: 1.5.1 + p-queue: 6.6.2 + ws: 8.21.1 + transitivePeerDependencies: + - bufferutil + - utf-8-validate + '@atomichub/atomicassets@2.0.2': {} '@bcoe/v8-coverage@1.0.2': {} @@ -3761,6 +3785,8 @@ snapshots: etag@1.8.1: {} + eventemitter3@4.0.7: {} + eventemitter3@5.0.4: {} events-universal@1.0.1: @@ -4433,6 +4459,8 @@ snapshots: dependencies: p-timeout: 6.1.4 + p-finally@1.0.0: {} + p-limit@3.1.0: dependencies: yocto-queue: 0.1.0 @@ -4441,11 +4469,20 @@ snapshots: dependencies: p-limit: 3.1.0 + p-queue@6.6.2: + dependencies: + eventemitter3: 4.0.7 + p-timeout: 3.2.0 + p-queue@9.3.1: dependencies: eventemitter3: 5.0.4 p-timeout: 7.0.1 + p-timeout@3.2.0: + dependencies: + p-finally: 1.0.0 + p-timeout@6.1.4: {} p-timeout@7.0.1: {} diff --git a/pnpm-workspace.yaml b/pnpm-workspace.yaml index 0512ac70..597bfe34 100644 --- a/pnpm-workspace.yaml +++ b/pnpm-workspace.yaml @@ -13,9 +13,10 @@ overrides: # The release-age gate exists to keep a freshly published third-party package # out of the lockfile before anyone has had a chance to notice it is hostile. -# This entry is first-party: the version is pinned, published from a tag in our -# own repository, and its tarball contents were checked against that tag. It can -# be dropped once the version is older than the gate; nothing depends on it -# staying. +# These entries are first-party: each version is pinned, published from a tag in +# our own repositories, and its tarball contents were checked against that tag. +# An entry can be dropped once the version is older than the gate; nothing +# depends on it staying. minimumReleaseAgeExclude: - '@atomichub/atomicassets@2.0.2' + - '@atomichub/antelope-ship-utils@1.0.0' diff --git a/src/bin/ship.ts b/src/bin/ship.ts deleted file mode 100644 index 03024eb0..00000000 --- a/src/bin/ship.ts +++ /dev/null @@ -1,25 +0,0 @@ -import StateHistoryBlockReader from '../connections/ship'; -import { ShipBlockResponse } from '../types/ship'; -import logger from '../utils/winston'; - -const ship = new StateHistoryBlockReader('ws://127.0.0.1:8080', { - min_block_confirmation: 1, - ds_threads: 1, - allow_empty_traces: true, - allow_empty_deltas: true, - allow_empty_blocks: true -}); - -ship.consume( (block: ShipBlockResponse) => { - logger.info('block received', block); - - ship.stopProcessing(); -}); - -ship.startProcessing({ - start_block_num: 97708771, - max_messages_in_flight: 1, - fetch_block: true, - fetch_traces: true, - fetch_deltas: true -}, ['contract_row']); diff --git a/src/connections/manager.ts b/src/connections/manager.ts index cda0ae14..3422c8b0 100644 --- a/src/connections/manager.ts +++ b/src/connections/manager.ts @@ -1,9 +1,10 @@ -import StateHistoryBlockReader from './ship'; +import { EOSJsDeserializer, StateHistoryConnection } from '@atomichub/antelope-ship-utils'; +import type { IShipConnectionOptions } from '@atomichub/antelope-ship-utils'; + import ChainApi from './chain'; import RedisConnection from './redis'; import PostgresConnection from './postgres'; import { IConnectionsConfig } from '../types/config'; -import { IBlockReaderOptions } from '../types/ship'; export default class ConnectionManager { readonly chain: ChainApi; @@ -96,13 +97,11 @@ export default class ConnectionManager { } } - createShipBlockReader(options?: IBlockReaderOptions): StateHistoryBlockReader { - const reader = new StateHistoryBlockReader(process.env.CHAIN_SHIP || this.config.chain.ship); - - if (options) { - reader.setOptions(options); - } - - return reader; + createShipConnection(connectionOptions: IShipConnectionOptions, deserializer: EOSJsDeserializer): StateHistoryConnection { + return new StateHistoryConnection({ + endpoint: process.env.CHAIN_SHIP || this.config.chain.ship, + connectionOptions, + deserializer + }); } } diff --git a/src/connections/ship.ts b/src/connections/ship.ts deleted file mode 100644 index 78199d97..00000000 --- a/src/connections/ship.ts +++ /dev/null @@ -1,438 +0,0 @@ -import PQueue from 'p-queue'; -import { ABI } from '@wharfkit/antelope'; -import WebSocket from 'ws'; -import { StaticPool } from 'node-worker-threads-pool'; - -import logger from '../utils/winston'; -import { - BlockRequestType, - IBlockReaderOptions, ShipBlockResponse -} from '../types/ship'; -import { deserializeEosioType, serializeEosioType } from '../utils/eosio'; - -export type BlockConsumer = (block: ShipBlockResponse) => any; - -export default class StateHistoryBlockReader { - currentArgs: BlockRequestType; - deltaWhitelist: string[]; - - abi: any; - shipAbi: ABI; - tables: Map; - - blocksQueue: PQueue; - - private ws: any; - - private connected: boolean; - private connecting: boolean; - private stopped: boolean; - - private deserializeWorkers: StaticPool<(x: Array<{type: string, data: Uint8Array, abi?: any}>) => any>; - - private unconfirmed: number; - private consumer: BlockConsumer; - - constructor( - private readonly endpoint: string, - private options: IBlockReaderOptions = {min_block_confirmation: 1, ds_threads: 4, allow_empty_deltas: false, allow_empty_traces: false, allow_empty_blocks: false} - ) { - this.connected = false; - this.connecting = false; - this.stopped = true; - - this.blocksQueue = new PQueue({concurrency: 1, autoStart: true}); - this.deserializeWorkers = undefined; - - this.consumer = null; - - this.abi = null; - this.shipAbi = null; - this.tables = new Map(); - - this.deltaWhitelist = []; - } - - setOptions(options?: Partial, deltas?: string[]): void { - if (options) { - this.options = {...this.options, ...options}; - } - - if (deltas) { - this.deltaWhitelist = deltas; - } - } - - connect(): void { - if (!this.connected && !this.connecting && !this.stopped) { - logger.info(`Connecting to ship endpoint ${this.endpoint}`); - logger.info(`Ship connect options ${JSON.stringify({...this.currentArgs, have_positions: 'removed'})}`); - - this.connecting = true; - - this.ws = new WebSocket(this.endpoint, { perMessageDeflate: false, maxPayload: 512 * 1024 * 1024 * 1024 }); - - this.ws.on('open', () => this.onConnect()); - this.ws.on('message', (data: any) => this.onMessage(data)); - this.ws.on('close', () => this.onClose()); - this.ws.on('error', (e: Error) => { logger.error('Websocket error', e); }); - } - } - - reconnect(): void { - if (this.stopped) { - return; - } - - setTimeout(() => { - logger.info('Reconnecting to Ship...'); - - this.connect(); - }, 5000); - } - - send(request: [string, any]): void { - this.ws.send(serializeEosioType('request', request, this.shipAbi)); - } - - onConnect(): void { - this.connected = true; - this.connecting = false; - } - - onMessage(data: any): void { - try { - if (!this.abi) { - logger.info('Receiving ABI from ship...'); - - this.abi = JSON.parse(data); - this.shipAbi = ABI.from(this.abi); - - if (this.options.ds_threads > 0) { - this.deserializeWorkers = new StaticPool({ - size: this.options.ds_threads, - task: './build/workers/deserializer.js', - workerData: {abi: JSON.stringify(this.abi)} - }); - } - - for (const table of this.abi.tables) { - this.tables.set(table.name, table.type); - } - - if (!this.stopped) { - this.requestBlocks(); - } - } else { - const [type, response] = deserializeEosioType('result', data, this.shipAbi); - - if (['get_blocks_result_v0', 'get_blocks_result_v1', 'get_blocks_result_v2'].indexOf(type) >= 0) { - const config: {[key: string]: {version: number }} = { - 'get_blocks_result_v0': {version: 0}, - 'get_blocks_result_v1': {version: 1}, - 'get_blocks_result_v2': {version: 2} - }; - - let block: any = null; - let traces: any = []; - let deltas: any = []; - - if (response.this_block) { - if (response.block) { - if (config[type].version === 2) { - block = this.deserializeParallel('signed_block_variant', response.block) - .then((res: any) => { - if (res[0] === 'signed_block_v1') { - return res[1]; - } - - throw new Error('Unsupported block type received ' + res[0]); - }); - } else if (config[type].version === 1) { - if (response.block[0] === 'signed_block_v1') { - block = response.block[1]; - } else { - block = Promise.reject(new Error('Unsupported block type received ' + response.block[0])); - } - } else if (config[type].version === 0) { - block = this.deserializeParallel('signed_block', response.block); - } else { - block = Promise.reject(new Error('Unsupported result type received ' + type)); - } - } else if(this.currentArgs.fetch_block) { - if (this.options.allow_empty_blocks) { - logger.warn('Block #' + response.this_block.block_num + ' does not contain block data'); - } else { - logger.error('Block #' + response.this_block.block_num + ' does not contain block data'); - - return this.blocksQueue.pause(); - } - } - - if (response.traces) { - traces = this.deserializeParallel('transaction_trace[]', response.traces); - } else if(this.currentArgs.fetch_traces) { - if (this.options.allow_empty_traces) { - logger.warn('Block #' + response.this_block.block_num + ' does not contain trace data'); - } else { - logger.error('Block #' + response.this_block.block_num + ' does not contain trace data'); - - return this.blocksQueue.pause(); - } - } - - if (response.deltas) { - deltas = this.deserializeParallel('table_delta[]', response.deltas) - .then(res => this.deserializeDeltas(res)); - } else if(this.currentArgs.fetch_deltas) { - if (this.options.allow_empty_deltas) { - logger.warn('Block #' + response.this_block.block_num + ' does not contain delta data'); - } else { - logger.error('Block #' + response.this_block.block_num + ' does not contain delta data'); - - return this.blocksQueue.pause(); - } - } - } - - this.blocksQueue.add(async () => { - if (response.this_block) { - this.currentArgs.start_block_num = response.this_block.block_num + 1; - } else { - this.currentArgs.start_block_num += 1; - } - - if (response.this_block && response.last_irreversible) { - this.currentArgs.have_positions = this.currentArgs.have_positions.filter( - row => row.block_num > response.last_irreversible.block_num && row.block_num < response.this_block.block_num - ); - - if (response.this_block.block_num > response.last_irreversible.block_num) { - this.currentArgs.have_positions.push(response.this_block); - } - } - - let deserializedTraces; - let deserializedDeltas; - - try { - deserializedTraces = await traces; - } catch (error) { - logger.error('Failed to deserialize traces at block #' + response.this_block.block_num, error); - - this.blocksQueue.clear(); - this.blocksQueue.pause(); - - throw error; - } - - try { - deserializedDeltas = await deltas; - } catch (error) { - logger.error('Failed to deserialize deltas at block #' + response.this_block.block_num, error); - - this.blocksQueue.clear(); - this.blocksQueue.pause(); - - throw error; - } - - try { - await this.processBlock({ - this_block: response.this_block, - head: response.head, - last_irreversible: response.last_irreversible, - prev_block: response.prev_block, - block: Object.assign( - {...response.this_block}, - await block, - {last_irreversible: response.last_irreversible}, - {head: response.head} - ), - traces: deserializedTraces, - deltas: deserializedDeltas - }); - } catch (error) { - logger.error('Ship blocks queue stopped due to an error at #' + response.this_block.block_num, error); - - this.blocksQueue.clear(); - this.blocksQueue.pause(); - - throw error; - } - - this.unconfirmed += 1; - - // Defer the ack while blocksQueue is overloaded so SHIP backs off. - // SHIP only sends past max_messages_in_flight when it receives an - // ack; holding the ack is the strongest backpressure we have on - // its send side. Accumulated `unconfirmed` is sent in one batch - // as soon as the queue drains under the threshold. See - // IBlockReaderOptions.max_blocks_queue for the full rationale. - const maxQueue = this.options.max_blocks_queue || 0; - const queueOverloaded = maxQueue > 0 && this.blocksQueue.size >= maxQueue; - - if (this.unconfirmed >= this.options.min_block_confirmation && !queueOverloaded) { - this.send(['get_blocks_ack_request_v0', { num_messages: this.unconfirmed }]); - this.unconfirmed = 0; - } - }).catch((error: any) => { - logger.error('Block processing error in ship queue', error); - }); - } else { - logger.warn('Not supported message received', {type, response}); - } - } - } catch (e) { - logger.error(e); - - this.ws.close(); - } - } - - async onClose(): Promise { - logger.error('Ship Websocket disconnected'); - - if (this.ws) { - await this.ws.terminate(); - this.ws = null; - } - - this.abi = null; - this.shipAbi = null; - this.tables = new Map(); - - this.connected = false; - this.connecting = false; - - this.blocksQueue.clear(); - - if (this.deserializeWorkers) { - await this.deserializeWorkers.destroy(); - this.deserializeWorkers = null; - } - - this.reconnect(); - } - - requestBlocks(): void { - this.unconfirmed = 0; - - this.send(['get_blocks_request_v0', this.currentArgs]); - } - - startProcessing(request: BlockRequestType = {}, deltas: string[] = []): void { - this.currentArgs = { - start_block_num: 0, - end_block_num: 0xffffffff, - max_messages_in_flight: 1, - have_positions: [], - irreversible_only: false, - fetch_block: false, - fetch_traces: false, - fetch_deltas: false, - ...request - }; - this.deltaWhitelist = deltas; - this.stopped = false; - - if (this.connected && this.abi) { - this.requestBlocks(); - } - - this.blocksQueue.start(); - - this.connect(); - } - - stopProcessing(): void { - this.stopped = true; - - this.ws.close(); - - this.blocksQueue.clear(); - this.blocksQueue.pause(); - } - - async processBlock(block: ShipBlockResponse): Promise { - if (!block.this_block) { - if (this.currentArgs.start_block_num >= this.currentArgs.end_block_num) { - logger.warn( - 'Empty block #' + this.currentArgs.start_block_num + ' received. Reader finished reading.' - ); - } else if (this.currentArgs.start_block_num % 10000 === 0) { - logger.warn( - 'Empty block #' + this.currentArgs.start_block_num + ' received. ' + - 'Node was likely started with a snapshot and you tried to process a block range ' + - 'before the snapshot. Catching up until init block.' - ); - } - - return; - } - - if (this.consumer) { - await this.consumer(block); - } - - return; - } - - consume(consumer: BlockConsumer): void { - this.consumer = consumer; - } - - private async deserializeParallel(type: string, data: Uint8Array): Promise { - if (this.options.ds_threads > 0) { - const result = await this.deserializeWorkers.exec([{type, data}]); - - if (result.success) { - return result.data[0]; - } - - throw new Error(result.message); - } - - return deserializeEosioType(type, data, this.shipAbi); - } - - private async deserializeArrayParallel(rows: Array<{type: string, data: Uint8Array}>): Promise { - if (this.options.ds_threads > 0) { - const result = await this.deserializeWorkers.exec(rows); - - if (result.success) { - return result.data; - } - - throw new Error(result.message); - } - - return rows.map(row => deserializeEosioType(row.type, row.data, this.shipAbi)); - } - - private async deserializeDeltas(deltas: any[]): Promise { - return await Promise.all(deltas.map(async (delta: any) => { - if (delta[0] === 'table_delta_v0' || delta[0] === 'table_delta_v1') { - if (this.deltaWhitelist.indexOf(delta[1].name) >= 0) { - const deserialized = await this.deserializeArrayParallel(delta[1].rows.map((row: any) => ({ - type: delta[1].name, data: row.data - }))); - - return [ - delta[0], - { - ...delta[1], - rows: delta[1].rows.map((row: any, index: number) => ({ - present: !!row.present, data: deserialized[index] - })) - } - ]; - } - - return delta; - } - - throw Error('Unsupported table delta type received ' + delta[0]); - })); - } -} diff --git a/src/filler/filler.ts b/src/filler/filler.ts index 810b8ef3..c445ba74 100644 --- a/src/filler/filler.ts +++ b/src/filler/filler.ts @@ -231,7 +231,7 @@ export default class Filler { lastBlockSpeeds.shift(); } - const queueState = `[DS:${this.reader.dsQueue.size}|SH:${this.reader.ship.blocksQueue.size}|JQ:${this.jobs.active}]`; + const queueState = `[DS:${this.reader.dsQueue.size}|SH:${this.reader.ship.getQueueSize()}|JQ:${this.jobs.active}]`; if (lastBlockNum === this.reader.currentBlock && lastBlockNum > 0) { const staleTime = Date.now() - lastBlockTime; diff --git a/src/filler/receiver-adapter.test.ts b/src/filler/receiver-adapter.test.ts new file mode 100644 index 00000000..1af1c78d --- /dev/null +++ b/src/filler/receiver-adapter.test.ts @@ -0,0 +1,218 @@ +import 'mocha'; +import {expect} from 'chai'; +import * as sinon from 'sinon'; +import { EventEmitter } from 'events'; +import { ShipError } from '@atomichub/antelope-ship-utils'; +import type { IBlockRequest, IShipConsumer } from '@atomichub/antelope-ship-utils'; + +import StateReceiver from './receiver'; +import ConnectionManager from '../connections/manager'; +import { IReaderConfig } from '../types/config'; +import { ModuleLoader } from './modules'; +import logger from '../utils/winston'; + +type ShipStub = EventEmitter & { + startProcessing: sinon.SinonStub; + stopProcessing: sinon.SinonStub; + getQueueSize: sinon.SinonStub; +}; + +type Handover = { + request?: IBlockRequest; + deltas?: string[]; +}; + +/** + * Build a StateReceiver through its real constructor against a stubbed + * ConnectionManager. The connection stub stands in for the package's + * StateHistoryConnection: it is an EventEmitter, and its startProcessing() + * pulls the request config and the required deltas off the consumer the way + * the package does, so the handover is exercised rather than asserted on a + * field. + */ +const sandbox = sinon.createSandbox(); + +function createReceiver(overrides: Partial = {}): { + receiver: StateReceiver; + ship: ShipStub; + createShipConnection: sinon.SinonStub; + handover: Handover; +} { + const ship = new EventEmitter() as ShipStub; + const handover: Handover = {}; + + ship.startProcessing = sandbox.stub().callsFake(async (consumer: IShipConsumer) => { + handover.request = await consumer.getRequestBlockConfig(); + handover.deltas = consumer.getRequiredDeltas(); + }); + ship.stopProcessing = sandbox.stub(); + ship.getQueueSize = sandbox.stub().returns(0); + + const createShipConnection = sandbox.stub().returns(ship); + + const connection = { + chain: { name: 'test-chain' }, + createShipConnection, + } as unknown as ConnectionManager; + + const config = { + name: 'test-reader', + start_block: 0, + stop_block: 0, + irreversible_only: false, + ship_prefetch_blocks: 50, + ship_min_block_confirmation: 10, + ship_ds_queue_size: 5, + ds_ship_threads: 0, + db_group_blocks: 12, + delete_data: false, + contracts: [], + ...overrides, + } as IReaderConfig; + + const receiver = new StateReceiver(config, connection, [], {} as ModuleLoader); + + return { receiver, ship, createShipConnection, handover }; +} + +function stubDatabase(receiver: StateReceiver, checkpoint: number, positions: any[]): void { + (receiver as any).database = { + getReaderPosition: sandbox.stub().resolves({ block_num: checkpoint, live: true, updated: 0 }), + getLastReaderBlocks: sandbox.stub().resolves(positions), + cleanupStaleReversibleData: sandbox.stub().resolves(), + }; +} + +describe('StateReceiver as an IShipConsumer', () => { + afterEach(() => { + sandbox.restore(); + }); + + describe('getRequestBlockConfig', () => { + it('carries the checkpoint start block and the positions the database returns', async () => { + const { receiver, handover } = createReceiver(); + const positions = [{ block_num: 5_999_995, block_id: 'a'.repeat(64) }]; + stubDatabase(receiver, 6_000_000, positions); + + await receiver.startProcessing(); + + expect(handover.request.start_block_num).to.equal(6_000_001); + expect(handover.request.have_positions).to.equal(positions); + expect(handover.request.max_messages_in_flight).to.equal(50); + expect(handover.request.end_block_num).to.equal(0xffffffff); + expect(handover.request.irreversible_only).to.equal(false); + expect(handover.request.fetch_block).to.equal(true); + expect(handover.request.fetch_traces).to.equal(true); + expect(handover.request.fetch_deltas).to.equal(true); + }); + + it('honours a configured stop block as the request end block', async () => { + const { receiver, handover } = createReceiver({ stop_block: 7_000_000 }); + stubDatabase(receiver, 6_000_000, []); + + await receiver.startProcessing(); + + expect(handover.request.end_block_num).to.equal(7_000_000); + }); + }); + + describe('prefetch and confirmation deadlock guard', () => { + it('lowers min_block_confirmation before the connection is constructed', () => { + const { createShipConnection } = createReceiver({ + ship_prefetch_blocks: 50, + ship_min_block_confirmation: 75, + }); + + expect(createShipConnection.calledOnce).to.equal(true); + // floor(50 / 2): SHIP would otherwise never receive a first ack, + // because the client waits for 75 processed blocks while the node + // stops sending after 50 unacked messages. + expect(createShipConnection.firstCall.args[0].min_block_confirmation).to.equal(25); + }); + + it('passes the configured confirmation count through when the prefetch depth covers it', () => { + const { createShipConnection } = createReceiver({ + ship_prefetch_blocks: 50, + ship_min_block_confirmation: 10, + }); + + expect(createShipConnection.firstCall.args[0].min_block_confirmation).to.equal(10); + }); + + it('never lowers the confirmation count below one', () => { + const { createShipConnection } = createReceiver({ + ship_prefetch_blocks: 1, + ship_min_block_confirmation: 5, + }); + + expect(createShipConnection.firstCall.args[0].min_block_confirmation).to.equal(1); + }); + + it('refuses empty payloads and forwards the queue ceiling to the connection', () => { + const { createShipConnection } = createReceiver({ ship_max_blocks_queue: 400 }); + const options = createShipConnection.firstCall.args[0]; + + expect(options.allow_empty_blocks).to.equal(false); + expect(options.allow_empty_traces).to.equal(false); + expect(options.allow_empty_deltas).to.equal(false); + expect(options.max_blocks_queue).to.equal(400); + }); + }); + + describe('connection events', () => { + it('reaches winston at the mapped level', () => { + // winston's leveled methods are overloaded, so the stubs are read + // back through the plain SinonStub shape to assert on call args. + const errorSpy = sandbox.stub(logger, 'error') as unknown as sinon.SinonStub; + const warnSpy = sandbox.stub(logger, 'warn') as unknown as sinon.SinonStub; + const infoSpy = sandbox.stub(logger, 'info') as unknown as sinon.SinonStub; + const debugSpy = sandbox.stub(logger, 'debug') as unknown as sinon.SinonStub; + + const { ship } = createReceiver(); + const shipError = new ShipError('Ship Websocket disconnected', new Error('ECONNRESET')); + + ship.emit('error', shipError); + ship.emit('warning', 'Block #5 does not contain delta data'); + ship.emit('info', 'Receiving ABI from ship...'); + ship.emit('debug', 'Block 5 processed'); + + expect(errorSpy.calledWith('Ship connection error', shipError)).to.equal(true); + expect(warnSpy.calledWith('Block #5 does not contain delta data')).to.equal(true); + expect(infoSpy.calledWith('Receiving ABI from ship...')).to.equal(true); + expect(debugSpy.calledWith('Block 5 processed')).to.equal(true); + }); + + it('carries the metadata a warning is emitted with', () => { + const warnSpy = sandbox.stub(logger, 'warn') as unknown as sinon.SinonStub; + const { ship } = createReceiver(); + const meta = { type: 'unknown_result_v9' }; + + ship.emit('warning', 'Not supported message received', meta); + + expect(warnSpy.calledWith('Not supported message received', meta)).to.equal(true); + }); + }); + + describe('getRequiredDeltas', () => { + it('asks for contract_row only', async () => { + const { receiver, handover } = createReceiver(); + stubDatabase(receiver, 6_000_000, []); + + expect(receiver.getRequiredDeltas()).to.deep.equal(['contract_row']); + + await receiver.startProcessing(); + + expect(handover.deltas).to.deep.equal(['contract_row']); + }); + }); + + describe('stopProcessing', () => { + it('stops the connection', async () => { + const { receiver, ship } = createReceiver(); + + await receiver.stopProcessing(); + + expect(ship.stopProcessing.calledOnce).to.equal(true); + }); + }); +}); diff --git a/src/filler/receiver.test.ts b/src/filler/receiver.test.ts index e0d6b734..2c330e98 100644 --- a/src/filler/receiver.test.ts +++ b/src/filler/receiver.test.ts @@ -11,7 +11,7 @@ import logger from '../utils/winston'; /** * Build a minimal StateReceiver-like object that has the fields used - * by the consumer() method, without requiring a real ConnectionManager. + * by the consume() method, without requiring a real ConnectionManager. */ function createReceiverStub(opts: { queueSize?: number; @@ -26,7 +26,7 @@ function createReceiverStub(opts: { const dsLock = new Semaphore(queueSize); const dsQueue = new PQueue({concurrency: 1, autoStart: true}); - // Build a partial StateReceiver with only the fields the consumer needs + // Build a partial StateReceiver with only the fields consume() needs const receiver = Object.create(StateReceiver.prototype) as StateReceiver; (receiver as any).dsLock = dsLock; (receiver as any).dsQueue = dsQueue; @@ -68,13 +68,12 @@ function makeBlockResponse(blockNum: number): ShipBlockResponse { } describe('StateReceiver', () => { - describe('consumer - dsLock semaphore management', () => { + describe('consume - dsLock semaphore management', () => { it('releases dsLock on successful block processing', async () => { const { receiver, dsLock } = createReceiverStub(); const resp = makeBlockResponse(1000); - // Call consumer (the private method) - await (receiver as any).consumer(resp); + await receiver.consume(resp); // Wait for dsQueue to drain await (receiver as any).dsQueue.onIdle(); @@ -96,7 +95,7 @@ describe('StateReceiver', () => { }); const resp = makeBlockResponse(2000); - await (receiver as any).consumer(resp); + await receiver.consume(resp); await dsQueue.onIdle(); // Even though process() threw, the dsLock should be released. @@ -117,7 +116,7 @@ describe('StateReceiver', () => { expect((receiver as any).queueStopped).to.equal(false); - await (receiver as any).consumer(makeBlockResponse(2500)); + await receiver.consume(makeBlockResponse(2500)); await dsQueue.onIdle(); // The fatal-stop branch flips queueStopped so filler.ts exits the pod @@ -130,7 +129,7 @@ describe('StateReceiver', () => { processResult: () => Promise.resolve(), }); - await (receiver as any).consumer(makeBlockResponse(2501)); + await receiver.consume(makeBlockResponse(2501)); await dsQueue.onIdle(); expect((receiver as any).queueStopped).to.equal(false); @@ -146,13 +145,13 @@ describe('StateReceiver', () => { }); // Process two blocks that will both fail - await (receiver as any).consumer(makeBlockResponse(3000)); + await receiver.consume(makeBlockResponse(3000)); await dsQueue.onIdle(); // Re-enable the queue (it gets paused on error) dsQueue.start(); - await (receiver as any).consumer(makeBlockResponse(3001)); + await receiver.consume(makeBlockResponse(3001)); await dsQueue.onIdle(); // With the fix: both permits are released, so we can acquire both @@ -174,9 +173,9 @@ describe('StateReceiver', () => { const { receiver, dsLock } = createReceiverStub({ prepareThrows: true }); const resp = makeBlockResponse(4000); - // consumer should throw (preprocessing failure), but dsLock must be released + // consume() should throw (preprocessing failure), but dsLock must be released try { - await (receiver as any).consumer(resp); + await receiver.consume(resp); } catch (_e) { // expected } @@ -199,7 +198,7 @@ describe('StateReceiver', () => { }); // Queue up multiple blocks - each acquires a dsLock permit - await (receiver as any).consumer(makeBlockResponse(5000)); + await receiver.consume(makeBlockResponse(5000)); // The first block will fail in the queue, triggering clear+purge. // After dsQueue drains, all permits should be recoverable. await dsQueue.onIdle(); @@ -215,86 +214,6 @@ describe('StateReceiver', () => { }); }); - describe('startProcessing - ACK deadlock prevention', () => { - it('overrides min_block_confirmation when prefetch < confirm would deadlock', async () => { - const dsLock = new Semaphore(5); - const dsQueue = new PQueue({concurrency: 1, autoStart: true}); - const setOptionsSpy = sinon.spy(); - - const receiver = Object.create(StateReceiver.prototype) as StateReceiver; - (receiver as any).dsLock = dsLock; - (receiver as any).dsQueue = dsQueue; - (receiver as any).config = { - name: 'test-reader', - ship_prefetch_blocks: 50, - ship_min_block_confirmation: 75, // BUG: 75 > 50 = deadlock - ship_ds_queue_size: 5, - start_block: 100, - stop_block: 0, - irreversible_only: false, - }; - (receiver as any).ship = { - setOptions: setOptionsSpy, - startProcessing: sinon.stub(), - consume: sinon.stub(), - }; - (receiver as any).database = { - getReaderPosition: sinon.stub().resolves({ block_num: 99, live: false, updated: 0 }), - getLastReaderBlocks: sinon.stub().resolves([]), - cleanupStaleReversibleData: sinon.stub().resolves(), - }; - (receiver as any).processor = { - setState: sinon.stub(), - }; - (receiver as any).handlers = []; - - await receiver.startProcessing(); - - // Should have called setOptions to override min_block_confirmation - expect(setOptionsSpy.calledOnce).to.equal(true); - const overrideOpts = setOptionsSpy.firstCall.args[0]; - expect(overrideOpts.min_block_confirmation).to.equal(25); // floor(50/2) - }); - - it('does not override when prefetch >= confirm', async () => { - const dsLock = new Semaphore(5); - const dsQueue = new PQueue({concurrency: 1, autoStart: true}); - const setOptionsSpy = sinon.spy(); - - const receiver = Object.create(StateReceiver.prototype) as StateReceiver; - (receiver as any).dsLock = dsLock; - (receiver as any).dsQueue = dsQueue; - (receiver as any).config = { - name: 'test-reader', - ship_prefetch_blocks: 50, - ship_min_block_confirmation: 10, // 10 < 50 = safe - ship_ds_queue_size: 5, - start_block: 100, - stop_block: 0, - irreversible_only: false, - }; - (receiver as any).ship = { - setOptions: setOptionsSpy, - startProcessing: sinon.stub(), - consume: sinon.stub(), - }; - (receiver as any).database = { - getReaderPosition: sinon.stub().resolves({ block_num: 99, live: false, updated: 0 }), - getLastReaderBlocks: sinon.stub().resolves([]), - cleanupStaleReversibleData: sinon.stub().resolves(), - }; - (receiver as any).processor = { - setState: sinon.stub(), - }; - (receiver as any).handlers = []; - - await receiver.startProcessing(); - - // Should NOT have called setOptions - expect(setOptionsSpy.called).to.equal(false); - }); - }); - describe('ship rewind guard', () => { function createStartedReceiver(checkpoint: number): StateReceiver { const receiver = Object.create(StateReceiver.prototype) as StateReceiver; @@ -310,9 +229,7 @@ describe('StateReceiver', () => { irreversible_only: false, }; (receiver as any).ship = { - setOptions: sinon.stub(), - startProcessing: sinon.stub(), - consume: sinon.stub(), + startProcessing: sinon.stub().resolves(), }; (receiver as any).database = { getReaderPosition: sinon.stub().resolves({ block_num: checkpoint, live: true, updated: 0 }), @@ -436,7 +353,7 @@ describe('StateReceiver', () => { (receiver as any).prepareContractRows = sinon.stub().resolves([]); (receiver as any).process = processSpy; - await (receiver as any).consumer(makeBlockResponse(9100)); + await receiver.consume(makeBlockResponse(9100)); await dsQueue.onIdle(); expect(processSpy.callCount).to.equal(1); diff --git a/src/filler/receiver.ts b/src/filler/receiver.ts index 7c4ffe55..0d6be900 100644 --- a/src/filler/receiver.ts +++ b/src/filler/receiver.ts @@ -1,9 +1,10 @@ import { ABI } from '@wharfkit/antelope'; import PQueue from 'p-queue'; +import { EOSJsDeserializer, StateHistoryConnection } from '@atomichub/antelope-ship-utils'; +import type { IBlockRequest, IShipConsumer } from '@atomichub/antelope-ship-utils'; import logger from '../utils/winston'; import ConnectionManager from '../connections/manager'; -import StateHistoryBlockReader from '../connections/ship'; import { IReaderConfig } from '../types/config'; import { ShipBlock, ShipBlockResponse, ShipTableDelta, ShipTransactionTrace } from '../types/ship'; import { EosioAction, EosioActionTrace, EosioContractRow, EosioTransaction } from '../types/eosio'; @@ -35,7 +36,7 @@ type ContractDataEstimation = { block_num: number }; -export default class StateReceiver { +export default class StateReceiver implements IShipConsumer { currentBlock = 0; headBlock = 0; lastIrreversibleBlock = 0; @@ -49,7 +50,7 @@ export default class StateReceiver { handlerDestructors: Array<() => void> = []; // Set true when the consumer queue stops on a non-recoverable block error - // (see consumer()). It's a known-dead state: the queue is cleared+paused and + // (see consume()). It's a known-dead state: the queue is cleared+paused and // no blocks will ever process again. The filler watchdog polls this and exits // immediately for a pod restart, instead of waiting out the multi-minute // "No blocks processed" stall timer (2026-06-01 incident - recovery took @@ -61,11 +62,18 @@ export default class StateReceiver { readonly dsLock: Semaphore; readonly dsQueue: PQueue; - readonly ship: StateHistoryBlockReader; + readonly ship: StateHistoryConnection; readonly processor: DataProcessor; readonly notifier: ApiNotificationSender; readonly database: ContractDB; + // The block request the package asks for through getRequestBlockConfig(). + // startProcessing() resolves the start block against the durable + // checkpoint before it hands this receiver to the connection, so the + // callback reads a settled value rather than recomputing one. + private startBlock = 0; + private readonly prefetch: number; + private readonly abis: {[key: string]: AbiCache}; constructor( @@ -84,19 +92,43 @@ export default class StateReceiver { this.notifier = new ApiNotificationSender(this.connection, this.processor, this.name); - this.ship = connection.createShipBlockReader({ - min_block_confirmation: config.ship_min_block_confirmation, - ds_threads: config.ds_ship_threads, + this.prefetch = config.ship_prefetch_blocks || 10; + + const minConfirm = config.ship_min_block_confirmation || 1; + let minBlockConfirmation = minConfirm; + + // SHIP stops sending once max_messages_in_flight messages are unacked, + // and the client acks only after min_block_confirmation processed + // blocks, so a prefetch depth below the confirmation count never + // reaches its first ack. The guard runs here, ahead of connection + // construction, rather than after a stall is observed. + if (this.prefetch < minConfirm) { + minBlockConfirmation = Math.max(1, Math.floor(this.prefetch / 2)); + + logger.warn( + `ship_prefetch_blocks (${this.prefetch}) < ship_min_block_confirmation (${minConfirm}) - ` + + `this will deadlock! Overriding min_block_confirmation to ${minBlockConfirmation}` + ); + } + + this.ship = connection.createShipConnection({ + min_block_confirmation: minBlockConfirmation, allow_empty_deltas: false, allow_empty_traces: false, allow_empty_blocks: false, - max_blocks_queue: config.ship_max_blocks_queue - }); + max_blocks_queue: config.ship_max_blocks_queue ?? 0 + }, new EOSJsDeserializer({ threads: config.ds_ship_threads })); + + // The connection reports through events. Without an 'error' listener an + // EventEmitter rethrows the emitted error, so this wiring is what keeps + // a socket failure a log line instead of a process-level throw. + this.ship.on('error', (error: Error) => logger.error('Ship connection error', error)); + this.ship.on('warning', (message: string, ...meta: any[]) => logger.warn(message, ...meta)); + this.ship.on('info', (message: string, ...meta: any[]) => logger.info(message, ...meta)); + this.ship.on('debug', (message: string, ...meta: any[]) => logger.debug(message, ...meta)); this.dsQueue = new PQueue({concurrency: 1, autoStart: true}); this.dsLock = new Semaphore(config.ship_ds_queue_size); - - this.ship.consume(this.consumer.bind(this)); } async startProcessing(): Promise { @@ -127,6 +159,8 @@ export default class StateReceiver { throw new Error('Reader end block cannot be lower than the starting block'); } + this.startBlock = startBlock; + logger.info('Reader ' + this.config.name + ' starting on block #' + startBlock); for (const handler of this.handlers) { @@ -150,28 +184,24 @@ export default class StateReceiver { // rows this stale survive when a past crash skipped the LIB-driven prune. await this.database.cleanupStaleReversibleData(this.name, this.lastIrreversibleBlock); - const prefetch = this.config.ship_prefetch_blocks || 10; - const minConfirm = this.config.ship_min_block_confirmation || 1; - - if (prefetch < minConfirm) { - const override = Math.max(1, Math.floor(prefetch / 2)); - logger.warn( - `ship_prefetch_blocks (${prefetch}) < ship_min_block_confirmation (${minConfirm}) - ` + - `this will deadlock! Overriding min_block_confirmation to ${override}` - ); - this.ship.setOptions({ min_block_confirmation: override }); - } + await this.ship.startProcessing(this); + } - this.ship.startProcessing({ - start_block_num: startBlock, + async getRequestBlockConfig(): Promise { + return { + start_block_num: this.startBlock, end_block_num: this.config.stop_block || 0xffffffff, - max_messages_in_flight: prefetch, + max_messages_in_flight: this.prefetch, irreversible_only: this.config.irreversible_only || false, have_positions: await this.database.getLastReaderBlocks(), fetch_block: true, fetch_traces: true, fetch_deltas: true - }, ['contract_row']); + }; + } + + getRequiredDeltas(): string[] { + return ['contract_row']; } async stopProcessing(): Promise { @@ -189,7 +219,7 @@ export default class StateReceiver { // without permitting a month-deep "fork" from a rewound SHIP node. private static readonly MAX_REVERSIBLE_WINDOW = 1000; - private async consumer(resp: ShipBlockResponse): Promise { + async consume(resp: ShipBlockResponse): Promise { await this.dsLock.acquire(); let actionTraces: Awaited>; diff --git a/src/types/ship.ts b/src/types/ship.ts index 5fedbbc3..e8fe7a24 100644 --- a/src/types/ship.ts +++ b/src/types/ship.ts @@ -1,150 +1,21 @@ -import { EosioContractRow } from './eosio'; - -export interface BlockRequestType { - start_block_num?: number; - end_block_num?: number; - max_messages_in_flight?: number; - have_positions?: any[]; - irreversible_only?: boolean; - fetch_block?: boolean; - fetch_traces?: boolean; - fetch_deltas?: boolean; -} - -export interface IBlockReaderOptions { - min_block_confirmation: number; - ds_threads: number; - allow_empty_traces: boolean; - allow_empty_deltas: boolean; - allow_empty_blocks: boolean; - /** - * Soft ceiling for the SHIP-side `blocksQueue`. When the queue is at or - * above this size, the worker stops acking blocks to SHIP, which makes - * SHIP back off (it stops sending past `max_messages_in_flight` without - * acks). Once the queue drains back under the threshold, accumulated - * `unconfirmed` is sent in one batch. Optional; when unset or 0, ack - * semantics are unchanged. Added 2026-05-09 (W2.2) after the WAX - * hype-drop cliff where blocksQueue grew to 994 - well past the - * 200-block max_messages_in_flight limit - and amplified every fork - * notification into a queue-wide rollback. - */ - max_blocks_queue?: number; -} - -export type ShipBlockResponse = { - head: {block_num: number, block_id: string}, - last_irreversible: {block_num: number, block_id: string}, - this_block: {block_num: number, block_id: string}, - prev_block: {block_num: number, block_id: string}, - block: ShipBlock, - traces: ShipTransactionTrace[], - deltas: ShipTableDelta[] -}; - -export type ShipBlock = { - block_num: number, - block_id: string, - head: {block_num: number, block_id: string}, - last_irreversible: {block_num: number, block_id: string}, - timestamp?: string, - producer?: string, - confirmed?: number, - previous?: string, - transaction_mroot?: string, - action_mroot?: string, - schedule_version?: number, - new_producers?: any | null, - header_extensions?: any[], - producer_signature?: string, - transactions?: any[], - block_extensions?: any[] -}; - -export type ShipTransactionTrace = [ - 'transaction_trace_v0', - { - id: string, - status: number, - cpu_usage_us: number, - net_usage_words: number, - elapsed: string, - net_usage: string, - scheduled: boolean, - action_traces: ShipActionTrace[], - account_ram_delta: Array<{account: string, delta: number}> | null, - except: any | null, - error_code: any | null, - failed_dtrx_trace: any | null, - partial: ShipPartialTransaction - } -]; - -export type ShipActionTrace = [ - 'action_trace_v0', - { - action_ordinal: number, - creator_action_ordinal: number, - receipt: ShipActionReceipt, - receiver: string, - act: { - account: string, - name: string, - authorization: Array<{actor: string, permission: string}>, - data: T - }, - context_free: boolean, - elapsed: string, - 'console': string, - account_ram_deltas: Array<{account: string, delta: number}>, - except: any | null, - error_code: any | null - } -]; - -export type ShipActionReceipt = [ - 'action_receipt_v0', - { - receiver: string, - act_digest: string, - global_sequence: string, - recv_sequence: string, - auth_sequence: Array<{account: string, sequence: string}>, - code_sequence: number, - abi_sequence: number - } -]; - -export type ShipPartialTransaction = [ - 'partial_transaction_v0', - { - expiration: string, - ref_block_num: number, - ref_block_prefix: number, - max_net_usage_words: number, - max_cpu_usage_ms: number, - delay_sec: number, - transaction_extensions: any[], - signatures: string[], - context_free_data: any[] - } -]; - -export type ShipTableDelta = [ - 'table_delta_v0', - { - name: string, - rows: Array<{present: boolean, data: [string, EosioContractRow]}> - } -]; - -export type ShipContractRow = [ - 'contract_row_v0', - { - code: string, - scope: string, - table: string, - primary_key: string, - payer: string, - value: T - } -]; +import type { + ShipActionTrace as PkgShipActionTrace, + ShipContractRow as PkgShipContractRow, + ShipTableDelta as PkgShipTableDelta +} from '@atomichub/antelope-ship-utils'; + +export type { + ShipBlock, + ShipBlockResponse, + ShipTransactionTrace, + ShipActionReceipt, + ShipPartialTransaction +} from '@atomichub/antelope-ship-utils'; + +// The package generics stay as published. This file pins the defaults this +// service declared for the trace and row aliases (binary action data is a +// hex string on some paths here) and re-exports ShipTableDelta without a +// type parameter. +export type ShipActionTrace = PkgShipActionTrace; +export type ShipContractRow = PkgShipContractRow; +export type ShipTableDelta = PkgShipTableDelta; diff --git a/src/workers/deserializer.ts b/src/workers/deserializer.ts deleted file mode 100644 index 14010864..00000000 --- a/src/workers/deserializer.ts +++ /dev/null @@ -1,35 +0,0 @@ -import { parentPort, workerData } from 'worker_threads'; -import { ABI } from '@wharfkit/antelope'; - -import logger from '../utils/winston'; -import { deserializeEosioType } from '../utils/eosio'; - -const args: {abi: string} = workerData; - -logger.info('Launching deserialization worker...'); - -const shipAbi: ABI = ABI.from(JSON.parse(args.abi)); - -parentPort.on('message', (param: Array<{type: string, data: Uint8Array | string, abi?: any}>) => { - try { - const result = []; - - for (const row of param) { - if (row.data === null) { - return parentPort.postMessage({success: false, message: 'Empty data received on deserialize worker'}); - } - - if (row.abi) { - const rowAbi = ABI.from(row.abi); - - result.push(deserializeEosioType(row.type, row.data, rowAbi)); - } else { - result.push(deserializeEosioType(row.type, row.data, shipAbi)); - } - } - - return parentPort.postMessage({success: true, data: result}); - } catch (e) { - return parentPort.postMessage({success: false, message: String(e)}); - } -});