diff --git a/.github/actions/local-network/action.yaml b/.github/actions/local-network/action.yaml index 142991c22df..68eb9849ee8 100644 --- a/.github/actions/local-network/action.yaml +++ b/.github/actions/local-network/action.yaml @@ -34,7 +34,7 @@ runs: ${{ env.HOME }}/.dashmate **/.env dashmate_volumes_dump - key: local-network-volumes/${{ steps.dashmate-fingerprint.outputs.sha }} + key: local-network-volumes/${{ steps.dashmate-fingerprint.outputs.sha }}/${{ hashFiles('.github/actions/local-network/action.yaml', 'packages/dapi/.env.example', 'packages/rs-drive-abci/.env.local', 'scripts/setup_local_network.sh', 'scripts/configure_test_suite.sh', 'scripts/configure_dotenv.sh', 'scripts/dashmate/volumes/**') }} - name: Restore dashmate volumes run: ./scripts/dashmate/volumes/restore.sh @@ -69,7 +69,7 @@ runs: ${{ env.HOME }}/.dashmate **/.env dashmate_volumes_dump - key: local-network-volumes/${{ steps.dashmate-fingerprint.outputs.sha }} + key: ${{ steps.local-network-data.outputs.cache-primary-key }} if: steps.local-network-data.outputs.cache-hit != 'true' - name: Configure pre-built docker images diff --git a/.github/workflows/tests-dashmate.yml b/.github/workflows/tests-dashmate.yml index 1dd01079edd..0f3cdbc38ab 100644 --- a/.github/workflows/tests-dashmate.yml +++ b/.github/workflows/tests-dashmate.yml @@ -18,7 +18,8 @@ jobs: dashmate-test: name: Run ${{ inputs.name }} tests runs-on: ubuntu-24.04 - timeout-minutes: 15 + # Cold helper-image pulls can consume most of 15 minutes before tests start. + timeout-minutes: 30 steps: - name: Check out repo uses: actions/checkout@v4 @@ -90,7 +91,7 @@ jobs: ${{ env.HOME }}/.dashmate **/.env dashmate_volumes_dump - key: local-network-volumes/${{ steps.dashmate-fingerprint.outputs.sha }} + key: local-network-volumes/${{ steps.dashmate-fingerprint.outputs.sha }}/${{ hashFiles('.github/actions/local-network/action.yaml', 'packages/dapi/.env.example', 'packages/rs-drive-abci/.env.local', 'scripts/setup_local_network.sh', 'scripts/configure_test_suite.sh', 'scripts/configure_dotenv.sh', 'scripts/dashmate/volumes/**') }} if: inputs.restore_local_network_data == true - name: Restore dashmate volumes @@ -110,7 +111,7 @@ jobs: env: DEBUG: 1 DASHMATE_E2E_TESTS_SKIP_IMAGE_BUILD: true - if: steps.local-network-data.outputs.cache-hit == 'false' + if: steps.local-network-data.outputs.cache-hit != 'true' - name: Show Docker logs if: ${{ failure() }} diff --git a/.github/workflows/tests-packges-functional.yml b/.github/workflows/tests-packges-functional.yml index e58fbccf210..59165b9d64c 100644 --- a/.github/workflows/tests-packges-functional.yml +++ b/.github/workflows/tests-packges-functional.yml @@ -5,7 +5,8 @@ jobs: test-functional: name: Run functional tests runs-on: ubuntu-24.04 - timeout-minutes: 15 + # Cold helper-image pulls can consume most of 15 minutes before tests start. + timeout-minutes: 30 env: ECR_HOST: ${{ secrets.AWS_ACCOUNT_ID }}.dkr.ecr.${{ vars.AWS_REGION }}.amazonaws.com steps: diff --git a/.github/workflows/tests-test-suite.yml b/.github/workflows/tests-test-suite.yml index 698e5ea153d..52ff20b9bde 100644 --- a/.github/workflows/tests-test-suite.yml +++ b/.github/workflows/tests-test-suite.yml @@ -23,7 +23,8 @@ jobs: test-suite: name: Run ${{ inputs.name }} runs-on: ubuntu-24.04 - timeout-minutes: 15 + # Cold helper-image pulls can consume most of 15 minutes before tests start. + timeout-minutes: 30 env: ECR_HOST: ${{ secrets.AWS_ACCOUNT_ID }}.dkr.ecr.${{ vars.AWS_REGION }}.amazonaws.com steps: diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml index 9ac8fa29957..fce0dd59e48 100644 --- a/.github/workflows/tests.yml +++ b/.github/workflows/tests.yml @@ -50,6 +50,7 @@ jobs: doctests-changed: ${{ steps.override.outputs.doctests-changed || steps.filter-doctests.outputs.doctests-changed }} swift-sdk-changed: ${{ steps.override.outputs.swift-sdk-changed || steps.filter-swift-sdk.outputs.swift-sdk-changed }} version-changed: ${{ steps.override.outputs.version-changed || steps.filter-version.outputs.version-changed }} + e2e-tests-changed: ${{ steps.override.outputs.e2e-tests-changed || steps.filter-e2e.outputs.e2e-tests-changed }} steps: - name: Checkout uses: actions/checkout@v4 @@ -89,6 +90,37 @@ jobs: - .github/workflows/tests.yml - .github/scripts/check-wallet-closure.py + - uses: dorny/paths-filter@v4 + id: filter-e2e + if: ${{ github.event_name != 'workflow_dispatch' }} + with: + filters: | + e2e-tests-changed: + - packages/platform-test-suite/** + - packages/dashmate/** + - packages/js-dapi-client/** + - packages/js-dash-sdk/** + - packages/wallet-lib/** + - packages/wasm-sdk/** + - packages/dapi/.env.example + - packages/rs-drive-abci/.env.local + - .github/actions/aws_ecr_login/** + - .github/actions/docker/** + - .github/actions/local-network/** + - .github/actions/nodejs/** + - .github/actions/rust/** + - .github/actions/sccache/** + - .github/workflows/tests.yml + - .github/workflows/tests-build-js.yml + - .github/workflows/tests-build-image.yml + - .github/workflows/tests-test-suite.yml + - .github/workflows/tests-packges-functional.yml + - .github/workflows/tests-dashmate.yml + - scripts/setup_local_network.sh + - scripts/configure_test_suite.sh + - scripts/configure_dotenv.sh + - scripts/dashmate/volumes/** + - name: Check for platform version change id: filter-version if: ${{ github.event_name != 'workflow_dispatch' }} @@ -187,6 +219,7 @@ jobs: with: filters: | swift-sdk-changed: + - .github/workflows/swift-sdk-build.yml - packages/swift-sdk/** - packages/dapi-grpc/** - packages/dashpay-contract/** @@ -265,12 +298,13 @@ jobs: echo 'doctests-changed=true' >> "$GITHUB_OUTPUT" echo 'swift-sdk-changed=true' >> "$GITHUB_OUTPUT" echo 'version-changed=true' >> "$GITHUB_OUTPUT" + echo 'e2e-tests-changed=true' >> "$GITHUB_OUTPUT" build-js: name: Build JS packages needs: - changes - if: ${{ needs.changes.outputs.js-packages != '[]' || needs.changes.outputs.version-changed == 'true' || github.event_name == 'schedule' || github.event_name == 'workflow_dispatch' }} + if: ${{ needs.changes.outputs.js-packages != '[]' || needs.changes.outputs.version-changed == 'true' || needs.changes.outputs.e2e-tests-changed == 'true' || github.event_name == 'schedule' || github.event_name == 'workflow_dispatch' }} secrets: inherit uses: ./.github/workflows/tests-build-js.yml @@ -279,10 +313,10 @@ jobs: needs: - check-secrets - changes - # Build Docker images on version change, nightly schedule, or manual dispatch + # Build Docker images for platform E2E changes, version changes, nightly schedules, or manual dispatches if: >- needs.check-secrets.outputs.has_ecr == 'true' && - (needs.changes.outputs.version-changed == 'true' || github.event_name == 'schedule' || github.event_name == 'workflow_dispatch') + (needs.changes.outputs.version-changed == 'true' || needs.changes.outputs.e2e-tests-changed == 'true' || github.event_name == 'schedule' || github.event_name == 'workflow_dispatch') secrets: inherit strategy: fail-fast: false @@ -374,8 +408,9 @@ jobs: uses: ./.github/workflows/tests-js-package.yml with: package: ${{ matrix.js-package }} - test-command: ${{ matrix.js-package == 'dashmate' && 'test:unit' || 'test' }} - skip-tests: ${{ contains(matrix.js-package, 'platform-test-suite') }} + # The platform test suite's default command drives a live network, so it + # runs from the E2E jobs. Its unit tests need nothing and run here. + test-command: ${{ (matrix.js-package == 'dashmate' || contains(matrix.js-package, 'platform-test-suite')) && 'test:unit' || 'test' }} direct-packages: ${{ needs.changes.outputs.js-packages-direct }} js-deps-versions: @@ -426,8 +461,9 @@ jobs: if: >- always() && needs.check-secrets.outputs.has_ecr == 'true' && + needs.build-js.result == 'success' && needs.build-images.result == 'success' && - (needs.changes.outputs.version-changed == 'true' || github.event_name == 'schedule' || github.event_name == 'workflow_dispatch') + (needs.changes.outputs.version-changed == 'true' || needs.changes.outputs.e2e-tests-changed == 'true' || github.event_name == 'schedule' || github.event_name == 'workflow_dispatch') test-suite: name: Test Suite @@ -442,7 +478,7 @@ jobs: needs.check-secrets.outputs.has_ecr == 'true' && needs.build-js.result == 'success' && needs.build-images.result == 'success' && - (needs.changes.outputs.version-changed == 'true' || github.event_name == 'schedule' || github.event_name == 'workflow_dispatch') + (needs.changes.outputs.version-changed == 'true' || needs.changes.outputs.e2e-tests-changed == 'true' || github.event_name == 'schedule' || github.event_name == 'workflow_dispatch') strategy: fail-fast: false matrix: @@ -479,5 +515,5 @@ jobs: needs.check-secrets.outputs.has_ecr == 'true' && needs.build-js.result == 'success' && needs.build-images.result == 'success' && - (needs.changes.outputs.version-changed == 'true' || github.event_name == 'schedule' || github.event_name == 'workflow_dispatch') + (needs.changes.outputs.version-changed == 'true' || needs.changes.outputs.e2e-tests-changed == 'true' || github.event_name == 'schedule' || github.event_name == 'workflow_dispatch') uses: ./.github/workflows/tests-packges-functional.yml diff --git a/packages/bench-suite/lib/client/createPlatformProofVerifier.js b/packages/bench-suite/lib/client/createPlatformProofVerifier.js index 420205435f7..aaa66e7fbf5 100644 --- a/packages/bench-suite/lib/client/createPlatformProofVerifier.js +++ b/packages/bench-suite/lib/client/createPlatformProofVerifier.js @@ -1,5 +1,59 @@ const DAPIAddress = require('@dashevo/dapi-client/lib/dapiAddressProvider/DAPIAddress'); +/** + * `WasmSdkError.name` for a transition family whose proof cannot bind the + * execution of one specific transition. + * + * @type {string} + */ +const EXECUTION_NOT_PROVED = 'ExecutionNotProved'; + +/** + * Read the error's kind without assuming it survives the WASM boundary. + * + * @param {*} error + * @returns {string|undefined} + */ +function readErrorName(error) { + try { + return error && error.name; + } catch (readError) { + // A WASM error object whose memory is already released throws on access. + return undefined; + } +} + +/** + * Convert a WASM SDK error into a plain `Error`. + * + * These are wasm-bindgen class instances rather than `Error`s, and they carry + * their kind and message on prototype getters. Mocha runs the suite in + * parallel workers and serializes a failure by copying the error's own + * properties, so an unconverted one arrives with nothing in it and the run + * reports a test that failed for no stated reason. + * + * @param {*} error + * @returns {Error} + */ +function toReportableError(error) { + if (error instanceof Error) { + return error; + } + + const name = readErrorName(error) || 'UnknownError'; + let message; + try { + message = (error && error.message) || String(error); + } catch (readError) { + message = 'error details are unavailable'; + } + + const reportable = new Error(`${name}: ${message}`); + reportable.name = name; + + return reportable; +} + /** * Shared EvoSDK instance. One per process: the verifier is stateless and the * underlying WASM SDK multiplexes concurrent requests. @@ -134,17 +188,25 @@ async function getEvoSdkForNetwork(callNetwork) { /** * Create an `IPlatformProofVerifier` for `Dash.Client`, backed by the - * Rust/WASM SDK, which authenticates every result end to end: GroveDB proof - * verification plus the Tenderdash quorum signature over the root hash. + * Rust/WASM SDK. Its proved re-query authenticates an execution result or + * affected-state snapshot end to end: GroveDB proof verification plus the + * Tenderdash quorum signature over the root hash. * * Verification re-queries Platform through the WASM SDK's proved paths rather - * than re-checking the exact bytes the JS transport received: the returned - * data (and the absence of a consensus error) is quorum-authenticated, so the - * unverified DAPI response is never the source of truth. + * than re-checking the exact bytes the JS transport received. Execution proof + * is required wherever the transition family can produce one. The families + * that cannot fall back to a height-pinned snapshot of the affected state, + * which is not evidence that this exact transition executed; the caller + * handles consensus errors from the original response before invoking this + * verifier. * + * @param {Object} [dependencies] + * @param {Function} [dependencies.getEvoSdkForNetwork] * @returns {Object} IPlatformProofVerifier */ -function createPlatformProofVerifier() { +function createPlatformProofVerifier({ + getEvoSdkForNetwork: loadEvoSdkForNetwork = getEvoSdkForNetwork, +} = {}) { return { /** * @param {Object} input @@ -153,16 +215,31 @@ function createPlatformProofVerifier() { * @returns {Promise} */ async verifyStateTransitionResult({ serializedStateTransition, network }) { - const { evo, sdk } = await getEvoSdkForNetwork(network); + const { evo, sdk } = await loadEvoSdkForNetwork(network); const stateTransition = evo.StateTransition.fromBytes( new Uint8Array(serializedStateTransition), ); - // Waits on the proved endpoint and verifies the execution proof and - // quorum signature inside the Rust SDK; throws unless the transition - // was executed (or yielded a consensus error, which also throws). - await sdk.stateTransitions.waitForResponse(stateTransition); + try { + await sdk.stateTransitions.waitForResponse(stateTransition); + } catch (error) { + // Balance top-ups, credit transfers and withdrawals, address funds + // movements, shields and no-history token operations have no proof + // that binds one specific transition, so the SDK reports that the + // execution was not proved rather than that anything failed. Their + // authenticated affected-state snapshot is the strongest result + // available; every other family keeps failing closed here. + if (readErrorName(error) !== EXECUTION_NOT_PROVED) { + throw toReportableError(error); + } + + try { + await sdk.stateTransitions.waitForAffectedState(stateTransition); + } catch (affectedStateError) { + throw toReportableError(affectedStateError); + } + } }, /** @@ -183,7 +260,7 @@ function createPlatformProofVerifier() { ); } - const { sdk } = await getEvoSdkForNetwork(network); + const { sdk } = await loadEvoSdkForNetwork(network); const history = await sdk.contracts.getHistory({ dataContractId: new Uint8Array(contractId), diff --git a/packages/dash-spv/lib/consensus.js b/packages/dash-spv/lib/consensus.js index 28c3963e8d7..8e9bbb54c55 100644 --- a/packages/dash-spv/lib/consensus.js +++ b/packages/dash-spv/lib/consensus.js @@ -65,6 +65,12 @@ function isValidBlockHeader(newHeader, previousHeaders, network = 'mainnet') { return false; } + // Dash Core disables difficulty retargeting on regtest. + if (network === 'regtest' && normalizedPreviousHeaders.length > 0) { + const previousHeader = normalizedPreviousHeaders[normalizedPreviousHeaders.length - 1]; + return normalizedHeader.bits === previousHeader.bits; + } + // A trusted checkpoint may contain fewer than a full DGW window. Proof of // work, the network pow limit, and timestamp rules still apply immediately; // exact difficulty validation begins as soon as the required history exists. diff --git a/packages/dash-spv/test/index.js b/packages/dash-spv/test/index.js index 1a15df18c13..00e843c3cd1 100644 --- a/packages/dash-spv/test/index.js +++ b/packages/dash-spv/test/index.js @@ -159,6 +159,27 @@ describe('Block header consensus validation', () => { consensus.isValidBlockHeader(nextHeader, wrongHistory, 'mainnet') .should.equal(false); }); + + it('should accept fast-mined regtest headers when Core retargeting is disabled', () => { + const previousHeaders = [new Blockchain('regtest').genesis]; + + while (previousHeaders.length < 24) { + previousHeaders.push(utils.createBlock(previousHeaders.at(-1), 0x207fffff)); + } + + const nextHeader = utils.createBlock(previousHeaders.at(-1), 0x207fffff); + + consensus.isValidBlockHeader(nextHeader, previousHeaders, 'regtest') + .should.equal(true); + }); + + it('should reject a changed regtest target before the DGW history window', () => { + const { genesis } = new Blockchain('regtest'); + const changedTargetHeader = utils.createBlock(genesis, 0x2070ffff); + + consensus.isValidBlockHeader(changedTargetHeader, [genesis], 'regtest') + .should.equal(false); + }); }); describe('SPV-DASH (addHeaders) add many headers for testnet', () => { diff --git a/packages/dashmate/src/createDIContainer.js b/packages/dashmate/src/createDIContainer.js index 876a2c00d4b..e8d98f4ddfb 100644 --- a/packages/dashmate/src/createDIContainer.js +++ b/packages/dashmate/src/createDIContainer.js @@ -47,6 +47,7 @@ import waitForConfirmations from './core/waitForConfirmations.js'; import generateBlsKeys from './core/generateBlsKeys.js'; import activateCoreSpork from './core/activateCoreSpork.js'; import waitForCorePeersConnected from './core/waitForCorePeersConnected.js'; +import waitForNodesToHaveTheSameHeight from './core/waitForNodesToHaveTheSameHeight.js'; import createNewAddress from './core/wallet/createNewAddress.js'; import generateBlocks from './core/wallet/generateBlocks.js'; @@ -256,6 +257,7 @@ export default async function createDIContainer(options = {}) { generateBlsKeys: asValue(generateBlsKeys), activateCoreSpork: asValue(activateCoreSpork), waitForCorePeersConnected: asValue(waitForCorePeersConnected), + waitForNodesToHaveTheSameHeight: asValue(waitForNodesToHaveTheSameHeight), }); /** diff --git a/packages/dashmate/src/listr/tasks/startGroupNodesTaskFactory.js b/packages/dashmate/src/listr/tasks/startGroupNodesTaskFactory.js index 53b6762ee1e..b35b3e252dc 100644 --- a/packages/dashmate/src/listr/tasks/startGroupNodesTaskFactory.js +++ b/packages/dashmate/src/listr/tasks/startGroupNodesTaskFactory.js @@ -4,12 +4,13 @@ import { NETWORK_LOCAL } from '../../constants.js'; import isServiceBuildRequired from '../../util/isServiceBuildRequired.js'; const { PrivateKey } = DashCoreLib; +const WAIT_FOR_NODES_TIMEOUT = 60 * 5 * 1000; /** * * @param {DockerCompose} dockerCompose * @param {waitForCorePeersConnected} waitForCorePeersConnected - * @param {waitForMasternodesSync} waitForMasternodesSync + * @param {waitForNodesToHaveTheSameHeight} waitForNodesToHaveTheSameHeight * @param {createRpcClient} createRpcClient * @param {Docker} docker * @param {startNodeTask} startNodeTask @@ -21,7 +22,7 @@ const { PrivateKey } = DashCoreLib; export default function startGroupNodesTaskFactory( dockerCompose, waitForCorePeersConnected, - waitForMasternodesSync, + waitForNodesToHaveTheSameHeight, createRpcClient, docker, startNodeTask, @@ -35,9 +36,14 @@ export default function startGroupNodesTaskFactory( * @return {Object} */ function startGroupNodesTask(configGroup) { + let coreRpcClients = []; + const minerConfig = configGroup.find((config) => ( config.get('core.miner.enable') )); + const isLocalMinerEnabled = () => ( + minerConfig && minerConfig.get('network') === NETWORK_LOCAL + ); const platformBuildConfig = configGroup.find((config) => ( isServiceBuildRequired(config) @@ -63,28 +69,36 @@ export default function startGroupNodesTaskFactory( }, { title: 'Wait for Core peers to be connected', - enabled: () => minerConfig && minerConfig.get('network') === NETWORK_LOCAL, - task: () => { - const tasks = configGroup.map((config) => ({ - title: `Checking ${config.getName()} peers`, - task: async () => { - const rpcClient = createRpcClient({ - port: config.get('core.rpc.port'), - user: 'dashmate', - pass: config.get('core.rpc.users.dashmate.password'), - host: await getConnectionHost(config, 'core', 'core.rpc.host'), - }); + enabled: isLocalMinerEnabled, + task: async () => { + coreRpcClients = await Promise.all(configGroup.map(async (config) => ( + createRpcClient({ + port: config.get('core.rpc.port'), + user: 'dashmate', + pass: config.get('core.rpc.users.dashmate.password'), + host: await getConnectionHost(config, 'core', 'core.rpc.host'), + }) + ))); - await waitForCorePeersConnected(rpcClient); - }, + const tasks = configGroup.map((config, index) => ({ + title: `Checking ${config.getName()} peers`, + task: () => waitForCorePeersConnected(coreRpcClients[index]), })); return new Listr(tasks, { concurrent: true }); }, }, + { + title: 'Wait for Core nodes to have the same height', + enabled: isLocalMinerEnabled, + task: () => waitForNodesToHaveTheSameHeight( + coreRpcClients, + WAIT_FOR_NODES_TIMEOUT, + ), + }, { title: 'Start a miner', - enabled: () => minerConfig && minerConfig.get('network') === NETWORK_LOCAL, + enabled: isLocalMinerEnabled, task: async () => { let minerAddress = minerConfig.get('core.miner.address'); diff --git a/packages/dashmate/test/unit/listr/tasks/startGroupNodesTaskFactory.spec.js b/packages/dashmate/test/unit/listr/tasks/startGroupNodesTaskFactory.spec.js new file mode 100644 index 00000000000..e83868ae9e0 --- /dev/null +++ b/packages/dashmate/test/unit/listr/tasks/startGroupNodesTaskFactory.spec.js @@ -0,0 +1,192 @@ +import startGroupNodesTaskFactory from '../../../../src/listr/tasks/startGroupNodesTaskFactory.js'; + +const WAIT_FOR_NODES_TIMEOUT = 60 * 5 * 1000; + +function createDeferred() { + let resolve; + const promise = new Promise((resolvePromise) => { + resolve = resolvePromise; + }); + + return { promise, resolve }; +} + +describe('startGroupNodesTaskFactory', () => { + function createConfig(sinon, name, rpcPort, network, minerEnabled = false) { + const values = { + 'core.miner.enable': minerEnabled, + 'core.miner.address': 'yN6Q6xj3Y9SuZ9pY4FAP8Zn46G7N8sCqzF', + 'core.miner.interval': 60, + 'core.rpc.port': rpcPort, + 'core.rpc.users.dashmate.password': `${name}-password`, + 'dashmate.helper.docker.build.enabled': false, + network, + 'platform.enable': false, + }; + + return { + get: sinon.stub().callsFake((path) => values[path]), + getName: sinon.stub().returns(name), + set: sinon.stub(), + }; + } + + function createFactory(sinon, overrides = {}) { + const dependencies = { + buildServicesTask: sinon.stub(), + createRpcClient: sinon.stub().callsFake((options) => ({ options })), + docker: {}, + dockerCompose: { + execCommand: sinon.stub().resolves(), + }, + getConnectionHost: sinon.stub().callsFake( + async (config) => `${config.getName()}.test`, + ), + startNodeTask: sinon.stub().resolves(), + waitForCorePeersConnected: sinon.stub().resolves(), + waitForNodeToBeReadyTask: sinon.stub(), + waitForNodesToHaveTheSameHeight: sinon.stub().resolves(), + ...overrides, + }; + + const startGroupNodesTask = startGroupNodesTaskFactory( + dependencies.dockerCompose, + dependencies.waitForCorePeersConnected, + dependencies.waitForNodesToHaveTheSameHeight, + dependencies.createRpcClient, + dependencies.docker, + dependencies.startNodeTask, + dependencies.waitForNodeToBeReadyTask, + dependencies.buildServicesTask, + dependencies.getConnectionHost, + ); + + return { dependencies, startGroupNodesTask }; + } + + it('should wait for every Core node to converge before starting the miner', async function shouldWaitForConvergence() { + const configs = [ + createConfig(this.sinon, 'local_seed', 19998, 'local', true), + createConfig(this.sinon, 'local_1', 20002, 'local'), + ]; + const convergence = createDeferred(); + const peerWaits = configs.map(() => createDeferred()); + const allPeersStarted = createDeferred(); + const firstRelevantCall = createDeferred(); + const peerClients = []; + + const { dependencies, startGroupNodesTask } = createFactory(this.sinon, { + dockerCompose: { + execCommand: this.sinon.stub().callsFake(async () => { + firstRelevantCall.resolve('miner'); + }), + }, + waitForCorePeersConnected: this.sinon.stub().callsFake((rpcClient) => { + peerClients.push(rpcClient); + if (peerClients.length === configs.length) { + allPeersStarted.resolve(); + } + + return peerWaits[peerClients.length - 1].promise; + }), + waitForNodesToHaveTheSameHeight: this.sinon.stub().callsFake(() => { + firstRelevantCall.resolve('convergence'); + return convergence.promise; + }), + }); + + const runPromise = startGroupNodesTask(configs).run({ + waitForReadiness: false, + }); + + try { + await allPeersStarted.promise; + expect(dependencies.waitForNodesToHaveTheSameHeight).to.not.have.been.called(); + expect(dependencies.dockerCompose.execCommand).to.not.have.been.called(); + + peerWaits[0].resolve(); + await Promise.resolve(); + expect(dependencies.waitForNodesToHaveTheSameHeight).to.not.have.been.called(); + + peerWaits[1].resolve(); + expect(await firstRelevantCall.promise).to.equal('convergence'); + expect(peerClients).to.have.lengthOf(configs.length); + expect(dependencies.dockerCompose.execCommand).to.not.have.been.called(); + + convergence.resolve(); + await runPromise; + + expect(dependencies.createRpcClient).to.have.callCount(configs.length); + expect(dependencies.createRpcClient.getCalls().map(({ args }) => args[0])) + .to.have.deep.members([ + { + port: 19998, + user: 'dashmate', + pass: 'local_seed-password', + host: 'local_seed.test', + }, + { + port: 20002, + user: 'dashmate', + pass: 'local_1-password', + host: 'local_1.test', + }, + ]); + + const rpcClients = dependencies.createRpcClient + .getCalls() + .map(({ returnValue }) => returnValue); + expect(peerClients).to.have.deep.members(rpcClients); + expect(dependencies.waitForNodesToHaveTheSameHeight) + .to.have.been.calledOnceWithExactly(rpcClients, WAIT_FOR_NODES_TIMEOUT); + expect(dependencies.dockerCompose.execCommand).to.have.been.calledOnce(); + } finally { + peerWaits.forEach(({ resolve }) => resolve()); + convergence.resolve(); + await runPromise.catch(() => {}); + } + }); + + it('should not start the miner when Core convergence fails', async function shouldStopOnFailure() { + const configs = [ + createConfig(this.sinon, 'local_seed', 19998, 'local', true), + createConfig(this.sinon, 'local_1', 20002, 'local'), + ]; + const convergenceError = new Error('Core tips do not match'); + const { dependencies, startGroupNodesTask } = createFactory(this.sinon, { + waitForNodesToHaveTheSameHeight: this.sinon.stub().rejects(convergenceError), + }); + + await expect(startGroupNodesTask(configs).run({ waitForReadiness: false })) + .to.be.rejectedWith(convergenceError); + + expect(dependencies.dockerCompose.execCommand).to.not.have.been.called(); + }); + + it('should not wait for Core convergence without a local miner', async function shouldSkipWithoutMiner() { + const configs = [ + createConfig(this.sinon, 'local_seed', 19998, 'local'), + createConfig(this.sinon, 'local_1', 20002, 'local'), + ]; + const { dependencies, startGroupNodesTask } = createFactory(this.sinon); + + await startGroupNodesTask(configs).run({ waitForReadiness: false }); + + expect(dependencies.createRpcClient).to.not.have.been.called(); + expect(dependencies.waitForNodesToHaveTheSameHeight).to.not.have.been.called(); + expect(dependencies.dockerCompose.execCommand).to.not.have.been.called(); + }); + + it('should not wait for Core convergence outside the local network', async function shouldSkipOutsideLocalNetwork() { + const configs = [ + createConfig(this.sinon, 'testnet', 19998, 'testnet', true), + ]; + const { dependencies, startGroupNodesTask } = createFactory(this.sinon); + + await startGroupNodesTask(configs).run({ waitForReadiness: false }); + + expect(dependencies.createRpcClient).to.not.have.been.called(); + expect(dependencies.waitForNodesToHaveTheSameHeight).to.not.have.been.called(); + expect(dependencies.dockerCompose.execCommand).to.not.have.been.called(); + }); +}); diff --git a/packages/js-dapi-client/lib/BlockHeadersProvider/BlockHeadersReader.js b/packages/js-dapi-client/lib/BlockHeadersProvider/BlockHeadersReader.js index f8d8fe4fc11..6b950f97650 100644 --- a/packages/js-dapi-client/lib/BlockHeadersProvider/BlockHeadersReader.js +++ b/packages/js-dapi-client/lib/BlockHeadersProvider/BlockHeadersReader.js @@ -8,6 +8,8 @@ const EVENTS = { ERROR: 'error', }; +const CLIENT_ALREADY_CLOSED_ERROR_MESSAGE = 'Client already closed - cannot .close()'; + /** * @typedef BlockHeadersReaderOptions * @property {number} [maxParallelStreams] @@ -132,8 +134,9 @@ class BlockHeadersReader extends EventEmitter { * @param e */ const rejectHeaders = async (e) => { - // Don't use cancelStream there because it's going to unsubscribe from events - stream.cancel(); + // Cancel the transport only: the retry below needs the listeners, so + // this cannot go through cancelStream, which unsubscribes from events + this.cancelStreamTransport(stream); stream.retryOnError(e); }; @@ -294,9 +297,25 @@ class BlockHeadersReader extends EventEmitter { } // eslint-disable-next-line class-methods-use-this + /** + * Cancels the stream's transport, tolerating one the transport has already + * closed. Listeners are left in place for callers that keep using the stream. + * @param {ReconnectableStream} stream + */ + // eslint-disable-next-line class-methods-use-this + cancelStreamTransport(stream) { + try { + stream.cancel(); + } catch (e) { + if (e.message !== CLIENT_ALREADY_CLOSED_ERROR_MESSAGE) { + throw e; + } + } + } + cancelStream(stream) { stream.removeAllListeners(); - stream.cancel(); + this.cancelStreamTransport(stream); } } diff --git a/packages/js-dapi-client/test/integration/BlockHeadersProvider/BlockHeadersProvider.spec.js b/packages/js-dapi-client/test/integration/BlockHeadersProvider/BlockHeadersProvider.spec.js index 1765aefbfcd..d92f1dc126d 100644 --- a/packages/js-dapi-client/test/integration/BlockHeadersProvider/BlockHeadersProvider.spec.js +++ b/packages/js-dapi-client/test/integration/BlockHeadersProvider/BlockHeadersProvider.spec.js @@ -166,4 +166,27 @@ describe('BlockHeadersProvider - integration', function describe() { expect(headHeight).to.equal(chainHeight + 1); await blockHeadersProvider.stop(); }); + + it('should accept the first historical batch after initializing from genesis', async function shouldAcceptFirstHistoricalBatch() { + const chainFromGenesis = await mockHeadersChain( + 'regtest', + 2, + undefined, + { mine: true }, + ); + await createBlockHeadersProvider(this.sinon, { + network: 'regtest', + targetBatchSize: historicalBatchSize, + }); + this.sinon.spy(blockHeadersProvider, 'emit'); + + await blockHeadersProvider.readHistorical(1, 1); + historicalStreams[0].sendHeaders(chainFromGenesis.slice(1)); + historicalStreams[0].end(); + + expect(blockHeadersProvider.spvChain.getLongestChain({ withPruned: true })) + .to.have.length(2); + expect(blockHeadersProvider.emit) + .to.have.been.calledWith(BlockHeadersProvider.EVENTS.HISTORICAL_DATA_OBTAINED); + }); }); diff --git a/packages/js-dapi-client/test/unit/BlockHeadersProvider/BlockHeadersReader.spec.js b/packages/js-dapi-client/test/unit/BlockHeadersProvider/BlockHeadersReader.spec.js index dca2784a9c0..63703479b95 100644 --- a/packages/js-dapi-client/test/unit/BlockHeadersProvider/BlockHeadersReader.spec.js +++ b/packages/js-dapi-client/test/unit/BlockHeadersProvider/BlockHeadersReader.spec.js @@ -137,6 +137,31 @@ describe('BlockHeadersReader - unit', () => { .to.have.been.calledTwice(); }); + it('should preserve a rejected headers error when the browser stream is already closed [data]', async function () { + blockHeadersReader.maxRetries = 0; + await blockHeadersReader.readHistorical(1, headers.length); + + const rejectWith = new Error('Invalid headers'); + const stream = historicalStreams[0]; + stream.cancel = this.sinon.stub(); + stream.cancel.onSecondCall() + .throws(new Error('Client already closed - cannot .close()')); + + blockHeadersReader.on(BlockHeadersReader.EVENTS.BLOCK_HEADERS, (_, rejectHeaders) => { + rejectHeaders(rejectWith); + }); + + const errorPromise = new Promise((resolve) => { + blockHeadersReader.on(BlockHeadersReader.EVENTS.ERROR, resolve); + }); + + expect(() => stream.sendHeaders(headers)).to.not.throw(); + expect(await errorPromise).to.equal(rejectWith); + expect(stream.cancel).to.have.been.calledTwice(); + expect(blockHeadersReader.stopReadingHistorical).to.have.been.calledOnce(); + expect(blockHeadersReader.historicalStreams).to.have.length(0); + }); + it('[error] should handle stream cancellation', async () => { await subscribeToHistoricalBatch(1, headers.length); diff --git a/packages/js-dash-sdk/src/SDK/Client/Platform/IPlatformProofVerifier.ts b/packages/js-dash-sdk/src/SDK/Client/Platform/IPlatformProofVerifier.ts index f5ed46e499d..bc8abf8871e 100644 --- a/packages/js-dash-sdk/src/SDK/Client/Platform/IPlatformProofVerifier.ts +++ b/packages/js-dash-sdk/src/SDK/Client/Platform/IPlatformProofVerifier.ts @@ -11,12 +11,17 @@ export interface VerifiedDataContractHistoryEntry { /** * Trust boundary for legacy JavaScript Platform operations. * - * Implementations must verify the GroveDB query/result and authenticate its - * root with the Tenderdash quorum signature for the supplied network and - * metadata. Returning successfully means the complete request binding was - * verified; presence or structural decoding of proof bytes is not sufficient. + * Implementations must verify GroveDB proofs and authenticate their roots with + * Tenderdash quorum signatures for the supplied network and metadata. + * Presence or structural decoding of proof bytes is not sufficient. */ export interface IPlatformProofVerifier { + /** + * Verify either the transition execution result or an authenticated, + * height-pinned snapshot of its affected state. A snapshot does not prove + * that this exact transition executed. Callers must reject consensus errors + * from the original response before invoking this method. + */ verifyStateTransitionResult(input: { serializedStateTransition: Uint8Array; response: IStateTransitionResult; @@ -24,6 +29,9 @@ export interface IPlatformProofVerifier { protocolVersion: number; }): Promise; + /** + * Verify the returned contract-history data and its complete query binding. + */ verifyDataContractHistory(input: { contractId: Uint8Array; startAtMs: bigint; diff --git a/packages/platform-test-suite/karma.conf.js b/packages/platform-test-suite/karma.conf.js index 05128836f99..37746721224 100644 --- a/packages/platform-test-suite/karma.conf.js +++ b/packages/platform-test-suite/karma.conf.js @@ -15,7 +15,10 @@ if (dotenvResult.error) { } // TODO: Fix test to be running in Browser -const testFilesPattern = './test/**/!(proofs|waitForStateTransitionResult).spec.js'; +// Only the suites that drive a network belong in the browser batches. Unit +// tests run under `test:unit`; sweeping them in here would let a unit failure +// abort a whole batch through `bail`. +const testFilesPattern = './test/@(e2e|functional)/**/!(proofs|waitForStateTransitionResult).spec.js'; const processors = ['webpack', 'sourcemap']; let testFiles = [ testFilesPattern, diff --git a/packages/platform-test-suite/lib/test/createPlatformProofVerifier.js b/packages/platform-test-suite/lib/test/createPlatformProofVerifier.js index 420205435f7..aaa66e7fbf5 100644 --- a/packages/platform-test-suite/lib/test/createPlatformProofVerifier.js +++ b/packages/platform-test-suite/lib/test/createPlatformProofVerifier.js @@ -1,5 +1,59 @@ const DAPIAddress = require('@dashevo/dapi-client/lib/dapiAddressProvider/DAPIAddress'); +/** + * `WasmSdkError.name` for a transition family whose proof cannot bind the + * execution of one specific transition. + * + * @type {string} + */ +const EXECUTION_NOT_PROVED = 'ExecutionNotProved'; + +/** + * Read the error's kind without assuming it survives the WASM boundary. + * + * @param {*} error + * @returns {string|undefined} + */ +function readErrorName(error) { + try { + return error && error.name; + } catch (readError) { + // A WASM error object whose memory is already released throws on access. + return undefined; + } +} + +/** + * Convert a WASM SDK error into a plain `Error`. + * + * These are wasm-bindgen class instances rather than `Error`s, and they carry + * their kind and message on prototype getters. Mocha runs the suite in + * parallel workers and serializes a failure by copying the error's own + * properties, so an unconverted one arrives with nothing in it and the run + * reports a test that failed for no stated reason. + * + * @param {*} error + * @returns {Error} + */ +function toReportableError(error) { + if (error instanceof Error) { + return error; + } + + const name = readErrorName(error) || 'UnknownError'; + let message; + try { + message = (error && error.message) || String(error); + } catch (readError) { + message = 'error details are unavailable'; + } + + const reportable = new Error(`${name}: ${message}`); + reportable.name = name; + + return reportable; +} + /** * Shared EvoSDK instance. One per process: the verifier is stateless and the * underlying WASM SDK multiplexes concurrent requests. @@ -134,17 +188,25 @@ async function getEvoSdkForNetwork(callNetwork) { /** * Create an `IPlatformProofVerifier` for `Dash.Client`, backed by the - * Rust/WASM SDK, which authenticates every result end to end: GroveDB proof - * verification plus the Tenderdash quorum signature over the root hash. + * Rust/WASM SDK. Its proved re-query authenticates an execution result or + * affected-state snapshot end to end: GroveDB proof verification plus the + * Tenderdash quorum signature over the root hash. * * Verification re-queries Platform through the WASM SDK's proved paths rather - * than re-checking the exact bytes the JS transport received: the returned - * data (and the absence of a consensus error) is quorum-authenticated, so the - * unverified DAPI response is never the source of truth. + * than re-checking the exact bytes the JS transport received. Execution proof + * is required wherever the transition family can produce one. The families + * that cannot fall back to a height-pinned snapshot of the affected state, + * which is not evidence that this exact transition executed; the caller + * handles consensus errors from the original response before invoking this + * verifier. * + * @param {Object} [dependencies] + * @param {Function} [dependencies.getEvoSdkForNetwork] * @returns {Object} IPlatformProofVerifier */ -function createPlatformProofVerifier() { +function createPlatformProofVerifier({ + getEvoSdkForNetwork: loadEvoSdkForNetwork = getEvoSdkForNetwork, +} = {}) { return { /** * @param {Object} input @@ -153,16 +215,31 @@ function createPlatformProofVerifier() { * @returns {Promise} */ async verifyStateTransitionResult({ serializedStateTransition, network }) { - const { evo, sdk } = await getEvoSdkForNetwork(network); + const { evo, sdk } = await loadEvoSdkForNetwork(network); const stateTransition = evo.StateTransition.fromBytes( new Uint8Array(serializedStateTransition), ); - // Waits on the proved endpoint and verifies the execution proof and - // quorum signature inside the Rust SDK; throws unless the transition - // was executed (or yielded a consensus error, which also throws). - await sdk.stateTransitions.waitForResponse(stateTransition); + try { + await sdk.stateTransitions.waitForResponse(stateTransition); + } catch (error) { + // Balance top-ups, credit transfers and withdrawals, address funds + // movements, shields and no-history token operations have no proof + // that binds one specific transition, so the SDK reports that the + // execution was not proved rather than that anything failed. Their + // authenticated affected-state snapshot is the strongest result + // available; every other family keeps failing closed here. + if (readErrorName(error) !== EXECUTION_NOT_PROVED) { + throw toReportableError(error); + } + + try { + await sdk.stateTransitions.waitForAffectedState(stateTransition); + } catch (affectedStateError) { + throw toReportableError(affectedStateError); + } + } }, /** @@ -183,7 +260,7 @@ function createPlatformProofVerifier() { ); } - const { sdk } = await getEvoSdkForNetwork(network); + const { sdk } = await loadEvoSdkForNetwork(network); const history = await sdk.contracts.getHistory({ dataContractId: new Uint8Array(contractId), diff --git a/packages/platform-test-suite/package.json b/packages/platform-test-suite/package.json index bb38e965b98..338e5557a63 100644 --- a/packages/platform-test-suite/package.json +++ b/packages/platform-test-suite/package.json @@ -6,6 +6,7 @@ "scripts": { "test": "yarn exec bin/test.sh", "lint": "eslint .", + "test:unit": "NODE_ENV=test mocha --no-config 'test/unit/**/*.spec.js'", "test:e2e": "NODE_ENV=test mocha 'test/e2e/**/*.spec.js'", "test:functional": "NODE_ENV=test mocha 'test/functional/**/*.spec.js'", "test:browsers": "karma start ./karma.conf.js" diff --git a/packages/platform-test-suite/test/e2e/dpns.spec.js b/packages/platform-test-suite/test/e2e/dpns.spec.js index 786f583f873..b9c6444682a 100644 --- a/packages/platform-test-suite/test/e2e/dpns.spec.js +++ b/packages/platform-test-suite/test/e2e/dpns.spec.js @@ -138,6 +138,11 @@ describe('DPNS', () => { const rawDocument = documents[0].toObject(); + // The registering identity is recorded as the creator. The local document + // returned by `register` never carries it, so the comparison below drops + // it on both sides and it is pinned here instead. + expect(Buffer.from(rawDocument.$creatorId)).to.deep.equal(identity.getId().toBuffer()); + delete rawDocument.$createdAt; delete rawDocument.$createdAtCoreBlockHeight; delete rawDocument.$createdAtBlockHeight; @@ -147,6 +152,7 @@ describe('DPNS', () => { delete rawDocument.$transferredAt; delete rawDocument.$transferredAtCoreBlockHeight; delete rawDocument.$transferredAtBlockHeight; + delete rawDocument.$creatorId; delete rawDocument.preorderSalt; const rawRegisteredDomain = registeredDomain.toObject(); @@ -160,6 +166,7 @@ describe('DPNS', () => { delete rawRegisteredDomain.$transferredAt; delete rawRegisteredDomain.$transferredAtCoreBlockHeight; delete rawRegisteredDomain.$transferredAtBlockHeight; + delete rawRegisteredDomain.$creatorId; delete rawRegisteredDomain.preorderSalt; expect(rawDocument).to.deep.equal(rawRegisteredDomain); @@ -179,6 +186,7 @@ describe('DPNS', () => { delete rawDocument.$transferredAt; delete rawDocument.$transferredAtCoreBlockHeight; delete rawDocument.$transferredAtBlockHeight; + delete rawDocument.$creatorId; delete rawDocument.preorderSalt; const rawRegisteredDomain = registeredDomain.toObject(); @@ -192,6 +200,7 @@ describe('DPNS', () => { delete rawRegisteredDomain.$transferredAt; delete rawRegisteredDomain.$transferredAtCoreBlockHeight; delete rawRegisteredDomain.$transferredAtBlockHeight; + delete rawRegisteredDomain.$creatorId; delete rawRegisteredDomain.preorderSalt; expect(rawDocument).to.deep.equal(rawRegisteredDomain); @@ -214,6 +223,7 @@ describe('DPNS', () => { delete rawDocument.$transferredAt; delete rawDocument.$transferredAtCoreBlockHeight; delete rawDocument.$transferredAtBlockHeight; + delete rawDocument.$creatorId; delete rawDocument.preorderSalt; const rawRegisteredDomain = registeredDomain.toObject(); @@ -227,6 +237,7 @@ describe('DPNS', () => { delete rawRegisteredDomain.$transferredAt; delete rawRegisteredDomain.$transferredAtCoreBlockHeight; delete rawRegisteredDomain.$transferredAtBlockHeight; + delete rawRegisteredDomain.$creatorId; delete rawRegisteredDomain.preorderSalt; expect(rawDocument).to.deep.equal(rawRegisteredDomain); diff --git a/packages/platform-test-suite/test/e2e/wallet.spec.js b/packages/platform-test-suite/test/e2e/wallet.spec.js index 241ba81204b..0a35ede23b1 100644 --- a/packages/platform-test-suite/test/e2e/wallet.spec.js +++ b/packages/platform-test-suite/test/e2e/wallet.spec.js @@ -4,8 +4,28 @@ const getDAPISeeds = require('../../lib/test/getDAPISeeds'); const createClientWithFundedWallet = require('../../lib/test/createClientWithFundedWallet'); const waitForBalanceToChange = require('../../lib/test/waitForBalanceToChange'); +const wait = require('../../lib/wait'); -const { EVENTS } = Dash.WalletLib; +const TRANSACTION_PROPAGATION_TIMEOUT_MS = 120000; +const TRANSACTION_POLL_INTERVAL_MS = 500; + +async function waitForTransaction(account, transactionId) { + const deadline = Date.now() + TRANSACTION_PROPAGATION_TIMEOUT_MS; + let transactions = account.getTransactions(); + + while (!transactions[transactionId] && Date.now() < deadline) { + await wait(TRANSACTION_POLL_INTERVAL_MS); + transactions = account.getTransactions(); + } + + if (!transactions[transactionId]) { + throw new Error( + `Transaction ${transactionId} did not reach the account within ${TRANSACTION_PROPAGATION_TIMEOUT_MS}ms`, + ); + } + + return transactions; +} describe('e2e', function e2eTest() { this.bail(true); @@ -107,15 +127,8 @@ describe('e2e', function e2eTest() { restoredAccount = await restoredWallet.getWalletAccount(); - let transactions = restoredAccount.getTransactions(); - - // Wait for new block if transaction has not been propagated yet - if (Object.keys(transactions).length === 0) { - await new Promise((resolve) => { restoredAccount.once(EVENTS.BLOCKHEADER, resolve); }); - transactions = restoredAccount.getTransactions(); - } - - await waitForBalanceToChange(restoredAccount); + // A mempool transaction may need the next block before a restored wallet can discover it. + const transactions = await waitForTransaction(restoredAccount, firstTransaction.id); const transactionIds = Object.keys(transactions); @@ -135,7 +148,8 @@ describe('e2e', function e2eTest() { waitForBalanceToChange(restoredAccount), ]); - const transactionIds = Object.keys(restoredAccount.getTransactions()); + const transactions = await waitForTransaction(restoredAccount, secondTransaction.id); + const transactionIds = Object.keys(transactions); expect(transactionIds).to.have.lengthOf(2); @@ -154,7 +168,8 @@ describe('e2e', function e2eTest() { await waitForBalanceToChange(emptyAccount); } - transactionIds = Object.keys(emptyAccount.getTransactions()); + const transactions = await waitForTransaction(emptyAccount, secondTransaction.id); + transactionIds = Object.keys(transactions); expect(transactionIds).to.have.lengthOf(2); diff --git a/packages/platform-test-suite/test/unit/createPlatformProofVerifier.spec.js b/packages/platform-test-suite/test/unit/createPlatformProofVerifier.spec.js new file mode 100644 index 00000000000..5794188d887 --- /dev/null +++ b/packages/platform-test-suite/test/unit/createPlatformProofVerifier.spec.js @@ -0,0 +1,136 @@ +const { expect, use } = require('chai'); +const chaiAsPromised = require('chai-as-promised'); +const sinon = require('sinon'); +const sinonChai = require('sinon-chai'); + +const createPlatformProofVerifier = require('../../lib/test/createPlatformProofVerifier'); + +use(chaiAsPromised); +use(sinonChai); + +/** + * Build a verifier over stubbed wait APIs. + * + * @param {Object} stateTransitions + * @returns {{ verifier: Object, evo: Object, stateTransition: Object }} + */ +function createVerifierWith(stateTransitions) { + const stateTransition = {}; + const evo = { + StateTransition: { + fromBytes: sinon.stub().returns(stateTransition), + }, + }; + const verifier = createPlatformProofVerifier({ + getEvoSdkForNetwork: sinon.stub().resolves({ evo, sdk: { stateTransitions } }), + }); + + return { verifier, evo, stateTransition }; +} + +/** + * The error the WASM SDK raises for a transition family that has no proof + * binding one specific transition. + * + * @returns {Error} + */ +function executionNotProvedError() { + const error = new Error('received a verified snapshot for this transition family'); + error.name = 'ExecutionNotProved'; + return error; +} + +describe('createPlatformProofVerifier', () => { + it('should require an execution proof before accepting a state transition', async () => { + const waitForResponse = sinon.stub().resolves(); + const waitForAffectedState = sinon.stub().resolves(); + const { verifier, evo, stateTransition } = createVerifierWith({ + waitForResponse, + waitForAffectedState, + }); + const serializedStateTransition = Uint8Array.from([1, 2, 3]); + + await verifier.verifyStateTransitionResult({ + serializedStateTransition, + network: 'local', + }); + + expect(evo.StateTransition.fromBytes).to.have.been.calledOnceWithExactly( + serializedStateTransition, + ); + expect(waitForResponse).to.have.been.calledOnceWithExactly(stateTransition); + expect(waitForAffectedState).to.not.have.been.called; + }); + + it('should accept an affected-state snapshot only when execution cannot be proved', async () => { + const waitForResponse = sinon.stub().rejects(executionNotProvedError()); + const waitForAffectedState = sinon.stub().resolves(); + const { verifier, stateTransition } = createVerifierWith({ + waitForResponse, + waitForAffectedState, + }); + + await verifier.verifyStateTransitionResult({ + serializedStateTransition: Uint8Array.from([1, 2, 3]), + network: 'local', + }); + + expect(waitForResponse).to.have.been.calledOnceWithExactly(stateTransition); + expect(waitForAffectedState).to.have.been.calledOnceWithExactly(stateTransition); + }); + + it('should propagate a failed execution proof without falling back to a snapshot', async () => { + const proofError = new Error('invalid execution proof'); + const waitForResponse = sinon.stub().rejects(proofError); + const waitForAffectedState = sinon.stub().resolves(); + const { verifier } = createVerifierWith({ waitForResponse, waitForAffectedState }); + + await expect(verifier.verifyStateTransitionResult({ + serializedStateTransition: Uint8Array.from([1, 2, 3]), + network: 'local', + })).to.be.rejectedWith(proofError); + + expect(waitForAffectedState).to.not.have.been.called; + }); + + it('should report a WASM error as a plain Error so its message survives', async () => { + // A wasm-bindgen error object: kind and message live on the prototype, and + // it is not an Error, so a parallel mocha worker serializes away everything + // it carries. + const wasmError = Object.create({ + get name() { + return 'Proof'; + }, + get message() { + return 'quorum signature is invalid'; + }, + }); + const { verifier } = createVerifierWith({ + waitForResponse: sinon.stub().rejects(wasmError), + waitForAffectedState: sinon.stub().resolves(), + }); + + const error = await verifier.verifyStateTransitionResult({ + serializedStateTransition: Uint8Array.from([1, 2, 3]), + network: 'local', + }).then(() => null, (thrown) => thrown); + + expect(error).to.be.an.instanceOf(Error); + expect(error.name).to.equal('Proof'); + expect(error.message).to.equal('Proof: quorum signature is invalid'); + expect(Object.getOwnPropertyNames(error)).to.include('message'); + }); + + it('should propagate affected-state proof rejection', async () => { + const proofError = new Error('invalid affected-state proof'); + const { verifier } = createVerifierWith({ + waitForResponse: sinon.stub().rejects(executionNotProvedError()), + waitForAffectedState: sinon.stub().rejects(proofError), + }); + + await expect(verifier.verifyStateTransitionResult({ + serializedStateTransition: Uint8Array.from([1, 2, 3]), + network: 'local', + })).to.be.rejectedWith(proofError); + }); +}); diff --git a/packages/rs-sdk-trusted-context-provider/src/provider.rs b/packages/rs-sdk-trusted-context-provider/src/provider.rs index f12f80610ce..4bf5cbe3c56 100644 --- a/packages/rs-sdk-trusted-context-provider/src/provider.rs +++ b/packages/rs-sdk-trusted-context-provider/src/provider.rs @@ -35,11 +35,13 @@ use std::error::Error as StdError; use std::net::ToSocketAddrs; use std::num::NonZeroUsize; use std::sync::{Arc, Mutex}; -#[cfg(not(target_arch = "wasm32"))] use std::time::Duration; use tracing::{debug, info}; use url::Url; +#[cfg(target_arch = "wasm32")] +const WASM_HTTP_REQUEST_TIMEOUT: Duration = Duration::from_secs(5); + /// A trusted HTTP-based context provider that fetches quorum information /// from trusted HTTP endpoints instead of requiring Core RPC access. #[derive(Clone)] @@ -90,6 +92,24 @@ struct MasternodeDiscoveryResponse { } impl TrustedHttpContextProvider { + /// Build a GET request for a trusted endpoint. + /// + /// Every request to `base_url` goes through here. On wasm32 the client + /// carries no timeout, because `ClientBuilder::timeout` is unavailable + /// there, so each request sets its own; a slow or hung endpoint would + /// otherwise block the caller forever. Cached responses are bypassed so a + /// rotated quorum or a changed masternode list is actually observed. + fn http_request(&self, url: &str) -> reqwest::RequestBuilder { + let request = self.client.get(url); + + #[cfg(target_arch = "wasm32")] + let request = request + .timeout(WASM_HTTP_REQUEST_TIMEOUT) + .fetch_cache_no_store(); + + request + } + /// Verify that a URL's domain resolves #[cfg(all( not(target_arch = "wasm32"), @@ -302,6 +322,33 @@ impl TrustedHttpContextProvider { Ok(()) } + /// Refresh current and previous quorum caches independently. + /// + /// Both endpoints are attempted so a failure from one does not prevent + /// usable data from the other from reaching the caches. + pub async fn refresh_quorum_caches(&self) -> Result<(), TrustedContextProviderError> { + let current_error = self.fetch_current_quorums().await.err(); + let previous_error = self.fetch_previous_quorums().await.err(); + + match (current_error, previous_error) { + (None, None) => Ok(()), + (Some(error), None) => Err(TrustedContextProviderError::NetworkError(format!( + "Failed to refresh current quorums: {}", + error + ))), + (None, Some(error)) => Err(TrustedContextProviderError::NetworkError(format!( + "Failed to refresh previous quorums: {}", + error + ))), + (Some(current_error), Some(previous_error)) => { + Err(TrustedContextProviderError::NetworkError(format!( + "Failed to refresh current quorums: {}; failed to refresh previous quorums: {}", + current_error, previous_error + ))) + } + } + } + /// Get the total number of quorums in both caches pub fn get_cached_quorum_count(&self) -> usize { let current_count = self @@ -329,7 +376,7 @@ impl TrustedHttpContextProvider { url ); - let response = self.client.get(&url).send().await?; + let response = self.http_request(&url).send().await?; if !response.status().is_success() { return Err(TrustedContextProviderError::NetworkError(format!( "HTTP {} from {}", @@ -392,7 +439,7 @@ impl TrustedHttpContextProvider { let url = format!("{}/quorums", self.base_url); debug!("Fetching current quorums from: {}", url); - let response = match self.client.get(&url).send().await { + let response = match self.http_request(&url).send().await { Ok(resp) => resp, Err(e) => { tracing::error!(error = ?e, url = %url, "HTTP request failed"); @@ -429,6 +476,11 @@ impl TrustedHttpContextProvider { debug!("Parsing JSON response for current quorums"); let quorums: QuorumsResponse = response.json().await?; + if !quorums.success { + return Err(TrustedContextProviderError::NetworkError( + "Current quorum response indicated failure".to_string(), + )); + } debug!("Successfully parsed {} quorums", quorums.data.len()); // Update cache @@ -465,7 +517,7 @@ impl TrustedHttpContextProvider { let url = format!("{}/previous", self.base_url); debug!("Fetching previous quorums from: {}", url); - let response = self.client.get(&url).send().await?; + let response = self.http_request(&url).send().await?; debug!("Received response with status: {}", response.status()); if !response.status().is_success() { @@ -478,6 +530,11 @@ impl TrustedHttpContextProvider { debug!("Parsing JSON response for previous quorums"); let quorums: PreviousQuorumsResponse = response.json().await?; + if !quorums.success { + return Err(TrustedContextProviderError::NetworkError( + "Previous quorum response indicated failure".to_string(), + )); + } debug!( "Successfully parsed {} previous quorums", quorums.data.quorums.len() @@ -836,6 +893,335 @@ impl ContextProvider for TrustedHttpContextProvider { #[cfg(test)] mod tests { use super::*; + use std::io::{BufRead, BufReader, Write}; + use std::net::{TcpListener, TcpStream}; + use std::thread; + use std::time::{Duration, Instant}; + + fn accept_before(listener: &TcpListener, deadline: Instant) -> TcpStream { + loop { + match listener.accept() { + Ok((stream, _)) => return stream, + Err(error) + if error.kind() == std::io::ErrorKind::WouldBlock + && Instant::now() < deadline => + { + thread::sleep(Duration::from_millis(10)); + } + Err(error) => panic!("accept quorum request: {}", error), + } + } + } + + fn current_response(hash: u8, key: u8) -> String { + serde_json::json!({ + "success": true, + "data": [{ + "quorum_hash": hex::encode([hash; 32]), + "key": hex::encode([key; 48]), + "height": 1, + "valid_members_count": 3 + }] + }) + .to_string() + } + + fn empty_current_response() -> String { + serde_json::json!({ + "success": true, + "data": [] + }) + .to_string() + } + + fn previous_response(hash: u8, key: u8) -> String { + serde_json::json!({ + "success": true, + "data": { + "height": 1, + "quorums": [{ + "quorum_hash": hex::encode([hash; 32]), + "key": hex::encode([key; 48]), + "height": 1, + "valid_members_count": 3 + }] + } + }) + .to_string() + } + + fn empty_previous_response() -> String { + serde_json::json!({ + "success": true, + "data": { + "height": 1, + "quorums": [] + } + }) + .to_string() + } + + fn spawn_http_responses( + responses: Vec<(&str, u16, String)>, + ) -> (String, thread::JoinHandle<()>) { + let listener = TcpListener::bind("127.0.0.1:0").expect("bind mock quorum endpoint"); + listener + .set_nonblocking(true) + .expect("make mock endpoint bounded"); + let address = listener.local_addr().expect("read mock endpoint address"); + let responses = responses + .into_iter() + .map(|(path, status, body)| (path.to_string(), status, body)) + .collect::>(); + + let handle = thread::spawn(move || { + for (expected_path, status, body) in responses { + let mut stream = accept_before(&listener, Instant::now() + Duration::from_secs(5)); + let mut reader = + BufReader::new(stream.try_clone().expect("clone quorum request stream")); + let mut request_line = String::new(); + reader + .read_line(&mut request_line) + .expect("read quorum request line"); + assert_eq!( + request_line.split_whitespace().nth(1), + Some(expected_path.as_str()) + ); + + loop { + let mut header = String::new(); + reader + .read_line(&mut header) + .expect("read quorum request header"); + if header == "\r\n" || header.is_empty() { + break; + } + } + + let reason = if status == 200 { + "OK" + } else { + "Internal Server Error" + }; + write!( + stream, + "HTTP/1.1 {} {}\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", + status, + reason, + body.len(), + body + ) + .expect("write quorum response"); + stream.flush().expect("flush quorum response"); + } + }); + + (format!("http://{}", address), handle) + } + + fn provider_for(base_url: String) -> TrustedHttpContextProvider { + TrustedHttpContextProvider::new_with_url( + Network::Regtest, + base_url, + NonZeroUsize::new(100).unwrap(), + ) + .expect("construct mock trusted context provider") + .with_refetch_if_not_found(false) + } + + #[tokio::test] + async fn refresh_quorum_caches_makes_rotated_keys_available() { + let (base_url, server) = spawn_http_responses(vec![ + ("/quorums", 200, current_response(0x11, 0x41)), + ("/previous", 200, previous_response(0x12, 0x42)), + ]); + let provider = provider_for(base_url); + + assert!(matches!( + provider + .get_quorum_public_key(1, [0x11; 32], 1) + .expect_err("rotated quorum must be absent before refresh"), + ContextProviderError::InvalidQuorum(_) + )); + provider + .refresh_quorum_caches() + .await + .expect("both quorum endpoints must refresh"); + assert_eq!( + provider + .get_quorum_public_key(1, [0x11; 32], 1) + .expect("current quorum must be cached"), + [0x41; 48] + ); + assert_eq!( + provider + .get_quorum_public_key(1, [0x12; 32], 1) + .expect("previous quorum must be cached"), + [0x42; 48] + ); + server.join().expect("mock quorum server must finish"); + } + + #[tokio::test] + async fn refresh_quorum_caches_attempts_both_endpoints_after_partial_failure() { + let (base_url, server) = spawn_http_responses(vec![ + ("/quorums", 500, "{}".to_string()), + ("/previous", 200, previous_response(0x22, 0x42)), + ]); + let provider = provider_for(base_url); + + assert!(provider.refresh_quorum_caches().await.is_err()); + assert_eq!( + provider + .get_quorum_public_key(1, [0x22; 32], 1) + .expect("previous quorum must be cached"), + [0x42; 48] + ); + server.join().expect("mock quorum server must finish"); + + let (base_url, server) = spawn_http_responses(vec![ + ("/quorums", 200, current_response(0x33, 0x43)), + ("/previous", 500, "{}".to_string()), + ]); + let provider = provider_for(base_url); + + assert!(provider.refresh_quorum_caches().await.is_err()); + assert_eq!( + provider + .get_quorum_public_key(1, [0x33; 32], 1) + .expect("current quorum must be cached"), + [0x43; 48] + ); + server.join().expect("mock quorum server must finish"); + } + + #[tokio::test] + async fn refresh_quorum_caches_preserves_cached_keys_when_both_endpoints_fail() { + let (base_url, server) = spawn_http_responses(vec![ + ("/quorums", 500, "{}".to_string()), + ("/previous", 500, "{}".to_string()), + ]); + let provider = provider_for(base_url); + provider.current_quorums_cache.lock().unwrap().put( + [0x44; 32], + QuorumData { + quorum_hash: hex::encode([0x44; 32]), + key: hex::encode([0x54; 48]), + height: 1, + valid_members_count: 3, + }, + ); + + assert!(provider.refresh_quorum_caches().await.is_err()); + assert_eq!( + provider + .get_quorum_public_key(1, [0x44; 32], 1) + .expect("known cached quorum must remain usable"), + [0x54; 48] + ); + assert!(matches!( + provider + .get_quorum_public_key(1, [0x45; 32], 1) + .expect_err("unknown quorum must remain rejected"), + ContextProviderError::InvalidQuorum(_) + )); + server.join().expect("mock quorum server must finish"); + } + + #[tokio::test] + async fn unsuccessful_quorum_response_does_not_populate_cache() { + let failed_current = serde_json::json!({ + "success": false, + "data": [{ + "quorum_hash": hex::encode([0x55; 32]), + "key": hex::encode([0x65; 48]), + "height": 1, + "valid_members_count": 3 + }] + }) + .to_string(); + let (base_url, server) = spawn_http_responses(vec![ + ("/quorums", 200, failed_current), + ("/previous", 200, empty_previous_response()), + ]); + let provider = provider_for(base_url); + + assert!(provider.refresh_quorum_caches().await.is_err()); + assert!(matches!( + provider + .get_quorum_public_key(1, [0x55; 32], 1) + .expect_err("unsuccessful response must not populate cache"), + ContextProviderError::InvalidQuorum(_) + )); + server.join().expect("mock quorum server must finish"); + + let failed_previous = serde_json::json!({ + "success": false, + "data": { + "height": 1, + "quorums": [{ + "quorum_hash": hex::encode([0x56; 32]), + "key": hex::encode([0x76; 48]), + "height": 1, + "valid_members_count": 3 + }] + } + }) + .to_string(); + let (base_url, server) = spawn_http_responses(vec![ + ("/quorums", 200, empty_current_response()), + ("/previous", 200, failed_previous), + ]); + let provider = provider_for(base_url); + provider.previous_quorums_cache.lock().unwrap().put( + [0x56; 32], + QuorumData { + quorum_hash: hex::encode([0x56; 32]), + key: hex::encode([0x66; 48]), + height: 1, + valid_members_count: 3, + }, + ); + + assert!(provider.refresh_quorum_caches().await.is_err()); + assert_eq!( + provider + .get_quorum_public_key(1, [0x56; 32], 1) + .expect("unsuccessful response must not overwrite cached quorum"), + [0x66; 48] + ); + server.join().expect("mock quorum server must finish"); + } + + #[tokio::test] + async fn successful_empty_refresh_preserves_cached_keys() { + let (base_url, server) = spawn_http_responses(vec![ + ("/quorums", 200, empty_current_response()), + ("/previous", 200, empty_previous_response()), + ]); + let provider = provider_for(base_url); + provider.current_quorums_cache.lock().unwrap().put( + [0x57; 32], + QuorumData { + quorum_hash: hex::encode([0x57; 32]), + key: hex::encode([0x67; 48]), + height: 1, + valid_members_count: 3, + }, + ); + + provider + .refresh_quorum_caches() + .await + .expect("successful empty responses are valid refreshes"); + assert_eq!( + provider + .get_quorum_public_key(1, [0x57; 32], 1) + .expect("empty refresh must not clear cached quorum"), + [0x67; 48] + ); + server.join().expect("mock quorum server must finish"); + } #[test] fn test_get_quorum_base_url() { diff --git a/packages/rs-sdk/src/platform/transition/put_document.rs b/packages/rs-sdk/src/platform/transition/put_document.rs index b34aea3512f..fa85a30a0dc 100644 --- a/packages/rs-sdk/src/platform/transition/put_document.rs +++ b/packages/rs-sdk/src/platform/transition/put_document.rs @@ -6,6 +6,7 @@ use crate::{Error, Sdk}; use dpp::dashcore::secp256k1::rand::rngs::StdRng; use dpp::dashcore::secp256k1::rand::{Rng, SeedableRng}; use dpp::data_contract::document_type::accessors::DocumentTypeV0Getters; +use dpp::data_contract::document_type::methods::DocumentTypeV0Methods; use dpp::data_contract::document_type::DocumentType; use dpp::document::{Document, DocumentV0Getters, DocumentV0Setters, INITIAL_REVISION}; use dpp::identity::signer::Signer; @@ -69,10 +70,11 @@ impl> PutDocument for Document { .await?; let settings = settings.unwrap_or_default(); + let document = prepare_document_for_transition(self, &document_type); let transition = if self.revision().is_some() && self.revision().unwrap() != INITIAL_REVISION { BatchTransition::new_document_replacement_transition_from_document( - self.clone(), + document, document_type.as_ref(), &identity_public_key, new_identity_contract_nonce, @@ -94,16 +96,16 @@ impl> PutDocument for Document { // before broadcasting to fail locally (no wasted nonce/fee). ensure_entropy_matches_document_id( &document_type.data_contract_id(), - &self.owner_id(), + &document.owner_id(), document_type.name(), &entropy, - self.id(), + document.id(), )?; - (self.clone(), entropy) + (document, entropy) } None => { let mut rng = StdRng::from_entropy(); - let mut document = self.clone(); + let mut document = document; let entropy = rng.gen::<[u8; 32]>(); document.set_id(Document::generate_document_id_v0( &document_type.data_contract_id(), @@ -161,6 +163,14 @@ impl> PutDocument for Document { } } +fn prepare_document_for_transition(document: &Document, document_type: &DocumentType) -> Document { + let mut document = document.clone(); + document_type + .as_ref() + .sanitize_document_properties(document.properties_mut()); + document +} + /// Ensures a caller-supplied `entropy` derives the same document id already set /// on a create document. /// @@ -197,6 +207,11 @@ fn ensure_entropy_matches_document_id( #[cfg(test)] mod tests { use super::*; + use dpp::data_contract::config::DataContractConfig; + use dpp::document::DocumentV0; + use dpp::platform_value::{platform_value, Value}; + use dpp::version::PlatformVersion; + use std::collections::BTreeMap; fn contract_id() -> Identifier { Identifier::from([1u8; 32]) @@ -252,4 +267,57 @@ mod tests { "a document id derived from a different entropy must be rejected locally" ); } + + #[test] + fn should_normalize_wasm_uint8_array_property_without_mutating_caller_document() { + let platform_version = PlatformVersion::latest(); + let config = DataContractConfig::default_for_version(platform_version) + .expect("should create default data contract config"); + let document_type = DocumentType::try_from_schema( + contract_id(), + 1, + config.version(), + "preorder", + platform_value!({ + "type": "object", + "properties": { + "saltedDomainHash": { + "type": "array", + "byteArray": true, + "minItems": 32_u32, + "maxItems": 32_u32, + "position": 0 + } + }, + "required": ["saltedDomainHash"], + "additionalProperties": false, + }), + None, + &BTreeMap::new(), + &config, + false, + &mut Vec::new(), + platform_version, + ) + .expect("should create DPNS-like document type"); + let integer_array = Value::Array(vec![Value::U64(7); 32]); + let document = Document::V0(DocumentV0 { + id: Identifier::new([3; 32]), + owner_id: owner_id(), + properties: BTreeMap::from([("saltedDomainHash".to_string(), integer_array.clone())]), + revision: Some(INITIAL_REVISION), + ..Default::default() + }); + + let prepared = prepare_document_for_transition(&document, &document_type); + + assert_eq!( + prepared.properties().get("saltedDomainHash"), + Some(&Value::Bytes32([7; 32])) + ); + assert_eq!( + document.properties().get("saltedDomainHash"), + Some(&integer_array) + ); + } } diff --git a/packages/swift-sdk/run_tests.sh b/packages/swift-sdk/run_tests.sh index 2df11b46dd1..47ca4b095d3 100755 --- a/packages/swift-sdk/run_tests.sh +++ b/packages/swift-sdk/run_tests.sh @@ -6,32 +6,120 @@ cd "$SCRIPT_DIR" || exit 1 # Provision an unlocked keychain for CI. # -# The `WalletStorage` tests write to the login keychain via `SecItemAdd`. +# The `WalletStorage` tests write to the user default keychain via `SecItemAdd`. # On a headless self-hosted runner (agent connected over SSH, no Aqua GUI # session) that call can fail with `errAuthorizationInternal` (-60008) # because the Security authorization subsystem has no session to service # the request, which surfaces as `keychainError(-60008)` and fails the # suite intermittently (it only passes when the runner happens to have a # live GUI session). +# Reusing a fixed CI keychain can also preserve stale permissions between +# jobs and surface as `errSecWrPerm` (-61). # -# Create a dedicated, explicitly-unlocked, no-auto-lock keychain and make -# it the default for the duration of the run so `SecItemAdd` targets a -# keychain that needs no interactive authorization. Gated to CI so it -# never touches a developer's login keychain, and the previous default is -# restored on exit so the persistent runner is left unchanged. +# Create a dedicated, explicitly-unlocked keychain and make it the user +# default for the duration of the run so `SecItemAdd` targets a keychain +# that needs no interactive authorization. Add it to the user search list +# so later reads and deletes find the same items. Gated to CI so it never +# touches a developer's keychain configuration; the previous default and +# search list are restored on exit. if [ -n "${CI:-}${GITHUB_ACTIONS:-}" ]; then - CI_KEYCHAIN="$HOME/Library/Keychains/dash-ci-tests.keychain-db" PREV_DEFAULT_KEYCHAIN="$(security default-keychain -d user | sed -E 's/^[[:space:]]*"?//;s/"?[[:space:]]*$//')" - restore_default_keychain() { - if [ -n "${PREV_DEFAULT_KEYCHAIN:-}" ] && [ -e "$PREV_DEFAULT_KEYCHAIN" ]; then - security default-keychain -d user -s "$PREV_DEFAULT_KEYCHAIN" || true + PREV_USER_KEYCHAINS_OUTPUT="$(security list-keychains -d user)" + PREV_USER_KEYCHAINS=() + while IFS= read -r keychain_path; do + keychain_path="$(printf '%s\n' "$keychain_path" | sed -E 's/^[[:space:]]*"?//;s/"?[[:space:]]*$//')" + if [ -n "$keychain_path" ]; then + PREV_USER_KEYCHAINS+=("$keychain_path") + fi + done <<< "$PREV_USER_KEYCHAINS_OUTPUT" + if [ "${#PREV_USER_KEYCHAINS[@]}" -eq 0 ]; then + echo "No user keychain search list is configured" >&2 + exit 1 + fi + + CI_KEYCHAIN_DIR="$(mktemp -d "${RUNNER_TEMP:-${TMPDIR:-/tmp}}/dash-ci-keychain.XXXXXX")" + CI_KEYCHAIN="$CI_KEYCHAIN_DIR/tests.keychain-db" + CI_KEYCHAIN_MAY_EXIST=0 + CI_DEFAULT_MAY_HAVE_CHANGED=0 + CI_SEARCH_LIST_MAY_HAVE_CHANGED=0 + + cleanup_ci_keychain() { + original_status=$? + cleanup_status=0 + trap - EXIT + + if [ "${CI_DEFAULT_MAY_HAVE_CHANGED:-0}" -eq 1 ]; then + if ! security default-keychain -d user -s "$PREV_DEFAULT_KEYCHAIN"; then + cleanup_status=1 + fi + fi + if [ "${CI_SEARCH_LIST_MAY_HAVE_CHANGED:-0}" -eq 1 ]; then + if ! security list-keychains -d user -s "${PREV_USER_KEYCHAINS[@]}"; then + cleanup_status=1 + fi + fi + if [ "${CI_KEYCHAIN_MAY_EXIST:-0}" -eq 1 ]; then + if ! security delete-keychain "$CI_KEYCHAIN"; then + cleanup_status=1 + fi + fi + if [ -d "${CI_KEYCHAIN_DIR:-}" ]; then + if ! rmdir "$CI_KEYCHAIN_DIR"; then + cleanup_status=1 + fi + fi + + if [ "$original_status" -ne 0 ]; then + exit "$original_status" + fi + if [ "$cleanup_status" -ne 0 ]; then + exit "$cleanup_status" fi } - trap restore_default_keychain EXIT - security create-keychain -p "" "$CI_KEYCHAIN" 2>/dev/null || true - security unlock-keychain -p "" "$CI_KEYCHAIN" - security set-keychain-settings "$CI_KEYCHAIN" # no auto-lock timeout + + trap cleanup_ci_keychain EXIT + CI_KEYCHAIN_DIR="$(cd "$CI_KEYCHAIN_DIR" && pwd -P)" + CI_KEYCHAIN="$CI_KEYCHAIN_DIR/tests.keychain-db" + chmod 700 "$CI_KEYCHAIN_DIR" + + CI_KEYCHAIN_PASSWORD="$(openssl rand -hex 32)" + CI_KEYCHAIN_MAY_EXIST=1 + security create-keychain -p "$CI_KEYCHAIN_PASSWORD" "$CI_KEYCHAIN" + security unlock-keychain -p "$CI_KEYCHAIN_PASSWORD" "$CI_KEYCHAIN" + security set-keychain-settings -u -t 7200 "$CI_KEYCHAIN" + CI_SEARCH_LIST_MAY_HAVE_CHANGED=1 + security list-keychains -d user -s "$CI_KEYCHAIN" "${PREV_USER_KEYCHAINS[@]}" + CI_DEFAULT_MAY_HAVE_CHANGED=1 security default-keychain -d user -s "$CI_KEYCHAIN" + + SELECTED_DEFAULT_KEYCHAIN="$(security default-keychain -d user | sed -E 's/^[[:space:]]*"?//;s/"?[[:space:]]*$//')" + if [ "$SELECTED_DEFAULT_KEYCHAIN" != "$CI_KEYCHAIN" ]; then + echo "Failed to select the CI test keychain as the user default" >&2 + exit 1 + fi + + KEYCHAIN_SMOKE_ACCOUNT="${GITHUB_ACTOR:-dash-ci}" + KEYCHAIN_SMOKE_SERVICE="dash-ci-keychain-${GITHUB_RUN_ID:-$$}-${GITHUB_RUN_ATTEMPT:-0}" + KEYCHAIN_SMOKE_VALUE="writable" + security add-generic-password \ + -a "$KEYCHAIN_SMOKE_ACCOUNT" \ + -s "$KEYCHAIN_SMOKE_SERVICE" \ + -w "$KEYCHAIN_SMOKE_VALUE" \ + "$CI_KEYCHAIN" + STORED_SMOKE_VALUE="$(security find-generic-password \ + -a "$KEYCHAIN_SMOKE_ACCOUNT" \ + -s "$KEYCHAIN_SMOKE_SERVICE" \ + -w \ + "$CI_KEYCHAIN")" + if [ "$STORED_SMOKE_VALUE" != "$KEYCHAIN_SMOKE_VALUE" ]; then + echo "CI test keychain read-back did not match the stored value" >&2 + exit 1 + fi + security delete-generic-password \ + -a "$KEYCHAIN_SMOKE_ACCOUNT" \ + -s "$KEYCHAIN_SMOKE_SERVICE" \ + "$CI_KEYCHAIN" + unset CI_KEYCHAIN_PASSWORD STORED_SMOKE_VALUE KEYCHAIN_SMOKE_VALUE fi # Pick a concrete iOS Simulator for the `xcodebuild test` run. A name diff --git a/packages/wallet-lib/src/CONSTANTS.js b/packages/wallet-lib/src/CONSTANTS.js index a367f1269cf..98f730499e8 100644 --- a/packages/wallet-lib/src/CONSTANTS.js +++ b/packages/wallet-lib/src/CONSTANTS.js @@ -101,7 +101,7 @@ const CONSTANTS = { UNKNOWN: 'unknown', }, STORAGE: { - version: 3, + version: 4, autosaveIntervalTime: 10 * 1000, REORG_SAFE_BLOCKS_COUNT: 6, }, diff --git a/packages/wallet-lib/src/plugins/Workers/BlockHeadersSyncWorker/BlockHeadersSyncWorker.js b/packages/wallet-lib/src/plugins/Workers/BlockHeadersSyncWorker/BlockHeadersSyncWorker.js index 76f7e02658c..c51d32be589 100644 --- a/packages/wallet-lib/src/plugins/Workers/BlockHeadersSyncWorker/BlockHeadersSyncWorker.js +++ b/packages/wallet-lib/src/plugins/Workers/BlockHeadersSyncWorker/BlockHeadersSyncWorker.js @@ -2,6 +2,7 @@ const BlockHeadersProvider = require('@dashevo/dapi-client/lib/BlockHeadersProvi const Worker = require('../../Worker'); const logger = require('../../../logger'); const EVENTS = require('../../../EVENTS'); +const deriveBlockHeadersResumeContext = require('../../../types/ChainStore/deriveBlockHeadersResumeContext'); const PROGRESS_UPDATE_INTERVAL = 1000; @@ -207,46 +208,32 @@ class BlockHeadersSyncWorker extends Worker { } /** - * Determines starting point considering options - * and last save checkpoint - * @returns {number|number} + * Determines the starting point from persisted authenticated headers. + * @returns {number} */ getStartBlockHeight() { const chainStore = this.storage.getDefaultChainStore(); - const bestBlockHeight = chainStore.state.chainHeight; - - let height; + const resumeContext = deriveBlockHeadersResumeContext(chainStore.state); const { - skipSynchronizationBeforeHeight, skipSynchronization, + skipSynchronizationBeforeHeight, } = (this.storage.application.syncOptions || {}); - if (skipSynchronization) { - this.logger.debug(`[BlockHeadersSyncWorker] Wallet created from a new mnemonic. Sync only last ${this.maxHeadersToKeep} blocks.`); - const syncFrom = bestBlockHeight - this.maxHeadersToKeep; - return syncFrom < 1 ? 1 : syncFrom; + // Headers are validated as a chain from an authenticated root, so a + // mid-chain start cannot be verified and these options cannot be honoured + // here. They still apply to transaction synchronization, so say so rather + // than let a caller assume header sync was skipped too. + if (skipSynchronization || skipSynchronizationBeforeHeight) { + this.logger.warn('[BlockHeadersSyncWorker] Header synchronization starts from the last authenticated header and ignores ' + + `${skipSynchronization ? 'skipSynchronization' : 'skipSynchronizationBeforeHeight'}; transaction synchronization still honours it`); } - const { lastSyncedHeaderHeight } = chainStore.state; - - if (typeof lastSyncedHeaderHeight !== 'number') { - throw new Error(`Invalid last synced header height ${lastSyncedHeaderHeight}`); - } - - const skipBefore = parseInt(skipSynchronizationBeforeHeight, 10); - - if (skipBefore > lastSyncedHeaderHeight) { - this.logger.debug(`[BlockHeadersSyncWorker] UNSAFE option skipSynchronizationBeforeHeight is set to ${skipBefore}`); - height = skipBefore; - } else if (lastSyncedHeaderHeight > -1) { - this.logger.debug(`[BlockHeadersSyncWorker] Last synced header height is ${lastSyncedHeaderHeight}`); - height = lastSyncedHeaderHeight; - } else { - height = 1; + if (resumeContext.startBlockHeight > 1) { + this.logger.debug(`[BlockHeadersSyncWorker] Last resumable header height is ${resumeContext.startBlockHeight}`); } - return height; + return resumeContext.startBlockHeight; } /** diff --git a/packages/wallet-lib/src/plugins/Workers/BlockHeadersSyncWorker/BlockHeadersSyncWorker.spec.js b/packages/wallet-lib/src/plugins/Workers/BlockHeadersSyncWorker/BlockHeadersSyncWorker.spec.js index 69f0da74cf6..e4ae35bc534 100644 --- a/packages/wallet-lib/src/plugins/Workers/BlockHeadersSyncWorker/BlockHeadersSyncWorker.spec.js +++ b/packages/wallet-lib/src/plugins/Workers/BlockHeadersSyncWorker/BlockHeadersSyncWorker.spec.js @@ -87,6 +87,9 @@ describe('BlockHeadersSyncWorker', () => { state: { chainHeight, lastSyncedHeaderHeight: -1, + blockHeaders: [], + headersMetadata: new Map(), + hashesByHeight: new Map(), }, updateLastSyncedHeaderHeight: sinon.spy(), updateChainHeight: sinon.spy(), @@ -112,7 +115,7 @@ describe('BlockHeadersSyncWorker', () => { expect(startBlockHeight).to.equal(1); }); - it('should return bestBlockHeight - N in case `skipSynchronization` option is present', () => { + it('should start at genesis for a new wallet without resumable headers', () => { /** * Mock options */ @@ -130,19 +133,41 @@ describe('BlockHeadersSyncWorker', () => { storage.getDefaultChainStore().state.chainHeight = 3000; startBlockHeight = blockHeadersSyncWorker.getStartBlockHeight(); - expect(startBlockHeight).to.equal(1000); + expect(startBlockHeight).to.equal(1); }); - it('should return last synced header height if present', () => { + it('should start at genesis for one stored header with a stale positive height', () => { const { storage } = blockHeadersSyncWorker; - storage.getDefaultChainStore().state.lastSyncedHeaderHeight = 1200; + const { state } = storage.getDefaultChainStore(); + state.blockHeaders = [spvChainHeaders[0]]; + state.lastSyncedHeaderHeight = 1200; + + const startBlockHeight = blockHeadersSyncWorker.getStartBlockHeight(); + + expect(startBlockHeight).to.equal(1); + }); + + it('should start at genesis for no stored headers with a stale positive height', () => { + const { state } = blockHeadersSyncWorker.storage.getDefaultChainStore(); + state.lastSyncedHeaderHeight = 1200; + + const startBlockHeight = blockHeadersSyncWorker.getStartBlockHeight(); + + expect(startBlockHeight).to.equal(1); + }); + + it('should return last synced header height for a resumable context', () => { + const { storage } = blockHeadersSyncWorker; + const { state } = storage.getDefaultChainStore(); + state.blockHeaders = spvChainHeaders.slice(0, 2); + state.lastSyncedHeaderHeight = 1200; const startBlockHeight = blockHeadersSyncWorker.getStartBlockHeight(); expect(startBlockHeight).to.equal(1200); }); - it('should return `skipSynchronizationBeforeHeight` value', () => { + it('should start at genesis instead of `skipSynchronizationBeforeHeight`', () => { const { storage } = blockHeadersSyncWorker; storage.application.syncOptions = { skipSynchronizationBeforeHeight: 1300, @@ -150,12 +175,14 @@ describe('BlockHeadersSyncWorker', () => { const startBlockHeight = blockHeadersSyncWorker.getStartBlockHeight(); - expect(startBlockHeight).to.equal(1300); + expect(startBlockHeight).to.equal(1); }); - it('should return last synced header if it\'s greater than `skipSynchronizationBeforeHeight` value', () => { + it('should return resumable header height when it is greater than the unsafe skip', () => { const { storage } = blockHeadersSyncWorker; - storage.getDefaultChainStore().state.lastSyncedHeaderHeight = 1300; + const { state } = storage.getDefaultChainStore(); + state.blockHeaders = spvChainHeaders.slice(0, 2); + state.lastSyncedHeaderHeight = 1300; storage.application.syncOptions = { skipSynchronizationBeforeHeight: 1200, }; @@ -165,16 +192,45 @@ describe('BlockHeadersSyncWorker', () => { expect(startBlockHeight).to.equal(1300); }); - it('should return `skipSynchronizationBeforeHeight` value if it\'s greater than last synced header height', () => { + it('should return resumable header height when the unsafe skip is greater', () => { const { storage } = blockHeadersSyncWorker; - storage.getDefaultChainStore().state.lastSyncedHeaderHeight = 1200; + const { state } = storage.getDefaultChainStore(); + state.blockHeaders = spvChainHeaders.slice(0, 2); + state.lastSyncedHeaderHeight = 1200; storage.application.syncOptions = { skipSynchronizationBeforeHeight: 1300, }; const startBlockHeight = blockHeadersSyncWorker.getStartBlockHeight(); - expect(startBlockHeight).to.equal(1300); + expect(startBlockHeight).to.equal(1200); + }); + + it('should start at genesis when the implied first header height is negative', () => { + const { state } = blockHeadersSyncWorker.storage.getDefaultChainStore(); + state.blockHeaders = spvChainHeaders.slice(0, 2); + state.lastSyncedHeaderHeight = 0; + + const startBlockHeight = blockHeadersSyncWorker.getStartBlockHeight(); + + expect(startBlockHeight).to.equal(1); + }); + + it('should reject invalid persisted header state types and ranges', () => { + const { state } = blockHeadersSyncWorker.storage.getDefaultChainStore(); + + state.blockHeaders = {}; + expect(() => blockHeadersSyncWorker.getStartBlockHeight()) + .to.throw('Invalid block headers'); + + state.blockHeaders = []; + state.lastSyncedHeaderHeight = Number.MAX_SAFE_INTEGER + 1; + expect(() => blockHeadersSyncWorker.getStartBlockHeight()) + .to.throw('Invalid last synced header height'); + + state.lastSyncedHeaderHeight = -2; + expect(() => blockHeadersSyncWorker.getStartBlockHeight()) + .to.throw('Invalid last synced header height'); }); }); @@ -252,7 +308,9 @@ describe('BlockHeadersSyncWorker', () => { it('should throw error if start block height is greater than best block height', async () => { const { storage } = blockHeadersSyncWorker; - storage.getDefaultChainStore().state.lastSyncedHeaderHeight = 2000; + const { state } = storage.getDefaultChainStore(); + state.blockHeaders = spvChainHeaders.slice(0, 2); + state.lastSyncedHeaderHeight = 2000; await expect(blockHeadersSyncWorker.onStart()) .to.be.rejectedWith('Start block height 2000 is greater than best block height 1000'); diff --git a/packages/wallet-lib/src/types/ChainStore/ChainStore.js b/packages/wallet-lib/src/types/ChainStore/ChainStore.js index 11788f8aff5..66e949bbd82 100644 --- a/packages/wallet-lib/src/types/ChainStore/ChainStore.js +++ b/packages/wallet-lib/src/types/ChainStore/ChainStore.js @@ -111,6 +111,15 @@ class ChainStore extends EventEmitter { this.state.lastSyncedHeaderHeight = height; } + resetBlockHeaders() { + Object.assign(this.state, { + blockHeaders: [], + lastSyncedHeaderHeight: -1, + headersMetadata: new Map(), + hashesByHeight: new Map(), + }); + } + updateLastSyncedBlockHeight(height) { if (height < this.state.lastSyncedBlockHeight) { throw new Error(`Cannot update lastSyncedBlockHeight to a lower value ${height} < ${this.state.lastSyncedBlockHeight}`); diff --git a/packages/wallet-lib/src/types/ChainStore/ChainStore.spec.js b/packages/wallet-lib/src/types/ChainStore/ChainStore.spec.js index 8892823c9ed..5563175568e 100644 --- a/packages/wallet-lib/src/types/ChainStore/ChainStore.spec.js +++ b/packages/wallet-lib/src/types/ChainStore/ChainStore.spec.js @@ -21,6 +21,45 @@ describe('ChainStore - class', () => { expect(testnetChainStore.state.instantLocks).to.deep.equal(new Map()); expect(testnetChainStore.state.addresses).to.deep.equal(new Map()); }); + + it('should reset only block header synchronization state', () => { + const chainStore = new ChainStore('testnet'); + const transactions = new Map([['transaction-id', { transaction: true }]]); + chainStore.state.chainHeight = 42; + chainStore.state.lastSyncedBlockHeight = 40; + chainStore.state.blockHeaders = [{ hash: 'header-hash' }]; + chainStore.state.lastSyncedHeaderHeight = 41; + chainStore.state.headersMetadata = new Map([['header-hash', { height: 41 }]]); + chainStore.state.hashesByHeight = new Map([[41, 'header-hash']]); + chainStore.state.transactions = transactions; + + chainStore.resetBlockHeaders(); + + expect(chainStore.state.blockHeaders).to.deep.equal([]); + expect(chainStore.state.lastSyncedHeaderHeight).to.equal(-1); + expect(chainStore.state.headersMetadata).to.deep.equal(new Map()); + expect(chainStore.state.hashesByHeight).to.deep.equal(new Map()); + expect(chainStore.state.chainHeight).to.equal(42); + expect(chainStore.state.lastSyncedBlockHeight).to.equal(40); + expect(chainStore.state.transactions).to.equal(transactions); + }); + + it('should preserve unsynced height sentinels when exporting a short chain', () => { + const chainStore = new ChainStore('regtest'); + chainStore.state.chainHeight = 0; + + const exportedState = chainStore.exportState(); + + expect(exportedState.lastSyncedHeaderHeight).to.equal(-1); + expect(exportedState.lastSyncedBlockHeight).to.equal(-1); + + const importedChainStore = new ChainStore('regtest'); + importedChainStore.importState(exportedState); + + expect(importedChainStore.state.lastSyncedHeaderHeight).to.equal(-1); + expect(importedChainStore.state.lastSyncedBlockHeight).to.equal(-1); + }); + it('should be able to import transactions with metadata', () => { const { transactions, txMetadata, lastSyncedBlockHeight, lastSyncedHeaderHeight, chainHeight } = fixtures1; diff --git a/packages/wallet-lib/src/types/ChainStore/deriveBlockHeadersResumeContext.js b/packages/wallet-lib/src/types/ChainStore/deriveBlockHeadersResumeContext.js new file mode 100644 index 00000000000..dc828c2e5b1 --- /dev/null +++ b/packages/wallet-lib/src/types/ChainStore/deriveBlockHeadersResumeContext.js @@ -0,0 +1,61 @@ +/** + * Derives the persisted header context that can safely resume synchronization. + * + * @param {ChainStoreState} state + * @returns {{ + * blockHeaders: Array, + * firstHeaderHeight: number, + * startBlockHeight: number, + * requiresHeaderStateReset: boolean, + * }} + */ +function deriveBlockHeadersResumeContext(state) { + const { + blockHeaders, + lastSyncedHeaderHeight, + headersMetadata, + hashesByHeight, + } = state; + + if (!Array.isArray(blockHeaders)) { + throw new Error('Invalid block headers: expected an array'); + } + + if (!Number.isSafeInteger(lastSyncedHeaderHeight) || lastSyncedHeaderHeight < -1) { + throw new Error(`Invalid last synced header height ${lastSyncedHeaderHeight}`); + } + + if (!(headersMetadata instanceof Map)) { + throw new Error('Invalid headers metadata: expected a Map'); + } + + if (!(hashesByHeight instanceof Map)) { + throw new Error('Invalid header hashes by height: expected a Map'); + } + + const firstHeaderHeight = lastSyncedHeaderHeight - blockHeaders.length + 1; + const canResume = blockHeaders.length >= 2 && firstHeaderHeight >= 0; + + if (canResume) { + return { + blockHeaders, + firstHeaderHeight, + startBlockHeight: lastSyncedHeaderHeight, + requiresHeaderStateReset: false, + }; + } + + const requiresHeaderStateReset = blockHeaders.length > 0 + || lastSyncedHeaderHeight !== -1 + || headersMetadata.size > 0 + || hashesByHeight.size > 0; + + return { + blockHeaders: [], + firstHeaderHeight: -1, + startBlockHeight: 1, + requiresHeaderStateReset, + }; +} + +module.exports = deriveBlockHeadersResumeContext; diff --git a/packages/wallet-lib/src/types/ChainStore/methods/exportState.js b/packages/wallet-lib/src/types/ChainStore/methods/exportState.js index 81235800693..89f9d8a6b42 100644 --- a/packages/wallet-lib/src/types/ChainStore/methods/exportState.js +++ b/packages/wallet-lib/src/types/ChainStore/methods/exportState.js @@ -20,7 +20,7 @@ function exportState() { fees: {}, }; - const reorgSafeHeight = chainHeight - STORAGE.REORG_SAFE_BLOCKS_COUNT; + const reorgSafeHeight = Math.max(-1, chainHeight - STORAGE.REORG_SAFE_BLOCKS_COUNT); const lastSyncedHeaderHeightToExport = Math.min(reorgSafeHeight, lastSyncedHeaderHeight); const lastSyncedBlockHeightToExport = Math.min(reorgSafeHeight, lastSyncedBlockHeight); diff --git a/packages/wallet-lib/src/types/Wallet/Wallet.js b/packages/wallet-lib/src/types/Wallet/Wallet.js index e6ed98aec62..741d1fb3ec1 100644 --- a/packages/wallet-lib/src/types/Wallet/Wallet.js +++ b/packages/wallet-lib/src/types/Wallet/Wallet.js @@ -192,6 +192,8 @@ class Wallet extends EventEmitter { } } + this.blockHeadersProviderInitializationPromise = null; + this.accountCreationPromises = new Map(); this.accounts = []; this.interface = opts.interface; // Suppressed global require to avoid cyclic dependencies diff --git a/packages/wallet-lib/src/types/Wallet/methods/createAccount.js b/packages/wallet-lib/src/types/Wallet/methods/createAccount.js index f74c4bb885c..c24e795a08e 100644 --- a/packages/wallet-lib/src/types/Wallet/methods/createAccount.js +++ b/packages/wallet-lib/src/types/Wallet/methods/createAccount.js @@ -1,5 +1,6 @@ const { WALLET_TYPES } = require('../../../CONSTANTS'); const EVENTS = require('../../../EVENTS'); +const deriveBlockHeadersResumeContext = require('../../ChainStore/deriveBlockHeadersResumeContext'); /** * Will derivate to a new account. * @param {object} accountOpts - options to pass, will autopopulate some @@ -27,6 +28,44 @@ async function createAccount(accountOpts) { if (this.walletType === WALLET_TYPES.SINGLE_ADDRESS) { baseOpts.privateKey = this.privateKey; } const opts = Object.assign(baseOpts, accountOpts); + if (!this.offlineMode) { + const chainStore = this.storage.getDefaultChainStore(); + const blockHeadersProvider = this.transport?.client?.blockHeadersProvider; + + if (blockHeadersProvider && !this.blockHeadersProviderInitializationPromise) { + this.blockHeadersProviderInitializationPromise = (async () => { + const resumeContext = deriveBlockHeadersResumeContext(chainStore.state); + + if (resumeContext.requiresHeaderStateReset) { + chainStore.resetBlockHeaders(); + + const hasPersistentAutosave = this.storage.autosave + && this.storage.adapter + && typeof this.storage.adapter.setItem === 'function'; + + if (hasPersistentAutosave) { + await this.storage.saveState(); + } + } + + await blockHeadersProvider.initializeChainWith( + resumeContext.blockHeaders, + resumeContext.firstHeaderHeight, + ); + })().catch((error) => { + // Forget a failed attempt so the next account creation retries. Keeping + // the rejected promise would replay the same failure for the lifetime + // of the wallet. + this.blockHeadersProviderInitializationPromise = null; + throw error; + }); + } + + if (this.blockHeadersProviderInitializationPromise) { + await this.blockHeadersProviderInitializationPromise; + } + } + const account = new Account(this, opts); // Add default derivation paths @@ -36,14 +75,6 @@ async function createAccount(accountOpts) { // at the moment of initialization (from persistent storage) account.createPathsForTransactions(); - // Add block headers from storage into the SPV chain if there are any - const chainStore = this.storage.getDefaultChainStore(); - const { blockHeaders, lastSyncedHeaderHeight } = chainStore.state; - if (!this.offlineMode && blockHeaders.length > 0) { - const { blockHeadersProvider } = this.transport.client; - blockHeadersProvider.initializeChainWith(blockHeaders, lastSyncedHeaderHeight); - } - this.accounts.push(account); try { diff --git a/packages/wallet-lib/src/types/Wallet/methods/createAccount.spec.js b/packages/wallet-lib/src/types/Wallet/methods/createAccount.spec.js index 7c41ff641ef..ee7f649e158 100644 --- a/packages/wallet-lib/src/types/Wallet/methods/createAccount.spec.js +++ b/packages/wallet-lib/src/types/Wallet/methods/createAccount.spec.js @@ -1,66 +1,492 @@ +const EventEmitter = require('events'); const { expect } = require('chai'); -const createAccount = require('./createAccount'); -const { WALLET_TYPES } = require('../../../CONSTANTS'); - -const exceptedException1 = 'getAccount expected index integer to be a property of accountOptions'; - - -// describe('Wallet - getAccount', () => { -// it('should warn on trying to pass arg as number', () => { -// let timesCreateAccountCalled = 0; -// let timesAttachEventsCalled = 0; -// const mockOpts = { -// accounts: [], -// storage: {}, -// walletType: WALLET_TYPES.HDWALLET, -// createAccount: (opts = { index: 0 }) => { -// timesCreateAccountCalled += 1; -// return { -// index: opts.index, -// storage: { -// attachEvents: () => timesAttachEventsCalled += 1, -// }, -// }; -// }, -// }; -// expect(() => getAccount.call(mockOpts, 0)).to.throw(exceptedException1); -// expect(timesCreateAccountCalled).to.equal(0); -// expect(timesAttachEventsCalled).to.equal(0); -// }); -// it('should create an account when not existing and get it back', () => { -// let timesCreateAccountCalled = 0; -// let timesAttachEventsCalled = 0; -// const mockOpts1 = { -// accounts: [], -// storage: {}, -// walletType: WALLET_TYPES.HDWALLET, -// createAccount: (opts = { index: 0 }) => { -// timesCreateAccountCalled += 1; -// const acc = { -// index: opts.index, -// storage: { -// attachEvents: () => timesAttachEventsCalled += 1, -// }, -// }; -// // This is actually done by Account class -// mockOpts1.accounts.push(acc); -// return acc; -// }, -// }; -// -// const acc = getAccount.call(mockOpts1); -// expect(acc.index).to.equal(0); -// expect(timesCreateAccountCalled).to.equal(1); -// expect(timesAttachEventsCalled).to.equal(1); -// const acc2 = getAccount.call(mockOpts1, { index: 0 }); -// expect(acc2.index).to.equal(0); -// expect(timesCreateAccountCalled).to.equal(1); -// expect(timesAttachEventsCalled).to.equal(2); -// expect(acc2).to.deep.equal(acc); -// -// const acc3 = getAccount.call(mockOpts1, { index: 1 }); -// expect(acc3.index).to.equal(1); -// expect(timesCreateAccountCalled).to.equal(2); -// expect(timesAttachEventsCalled).to.equal(3); -// }); -// }); +const DAPIClient = require('@dashevo/dapi-client'); + +const Wallet = require('../Wallet'); +const EVENTS = require('../../../EVENTS'); +const logger = require('../../../logger'); +const BlockHeadersSyncWorker = require('../../../plugins/Workers/BlockHeadersSyncWorker/BlockHeadersSyncWorker'); +const { mineHeadersChain } = require('../../../test/mocks/dashcore/block'); + +const { BlockHeadersProvider } = DAPIClient; + +class TestTransport { + constructor(blockHeadersProvider) { + this.client = { + blockHeadersProvider, + }; + } +} + +class TestTransportWithoutClient {} + +const waitForStorage = async (wallet) => { + if (!wallet.storage.configured) { + await new Promise((resolve) => { + wallet.storage.once(EVENTS.CONFIGURED, resolve); + }); + } +}; + +describe('Wallet#createAccount', function suite() { + let wallet; + + afterEach(async () => { + if (wallet) { + await wallet.disconnect(); + } + }); + + it('should wait for genesis initialization when the online store has no headers', async function shouldInitializeEmptyStore() { + let resolveInitialization; + const initializationPromise = new Promise((resolve) => { + resolveInitialization = resolve; + }); + const blockHeadersProvider = { + initializeChainWith: this.sinon.stub().returns(initializationPromise), + }; + wallet = new Wallet({ + mnemonic: null, + network: 'regtest', + transport: new TestTransport(blockHeadersProvider), + }); + + let accountCreated = false; + const createAccountPromise = wallet.createAccount({ + injectDefaultPlugins: false, + synchronize: false, + }).then((account) => { + accountCreated = true; + return account; + }); + + await new Promise((resolve) => { + setTimeout(resolve, 0); + }); + + expect(blockHeadersProvider.initializeChainWith) + .to.have.been.calledOnceWithExactly([], -1); + expect(accountCreated).to.equal(false); + + resolveInitialization(); + await createAccountPromise; + }); + + it('should resume from the first authenticated stored header without a remote checkpoint', async function shouldInitializeStoredHeaders() { + const blockHeadersProvider = { + initializeChainWith: this.sinon.stub().resolves(), + }; + wallet = new Wallet({ + mnemonic: null, + network: 'regtest', + transport: new TestTransport(blockHeadersProvider), + }); + await waitForStorage(wallet); + const chainStore = wallet.storage.getDefaultChainStore(); + const storedHeaders = [ + { hash: 'stored-header-40' }, + { hash: 'stored-header-41' }, + { hash: 'stored-header-42' }, + ]; + chainStore.state.blockHeaders = storedHeaders; + chainStore.state.lastSyncedHeaderHeight = 42; + + await wallet.createAccount({ + injectDefaultPlugins: false, + synchronize: false, + }); + + expect(blockHeadersProvider.initializeChainWith) + .to.have.been.calledOnceWithExactly(storedHeaders, 40); + }); + + it('should normalize a stale one-header context before genesis synchronization', async function shouldNormalizeInsufficientHeaders() { + const blockHeadersProvider = new BlockHeadersProvider({ network: 'regtest' }); + wallet = new Wallet({ + mnemonic: null, + network: 'regtest', + transport: new TestTransport(blockHeadersProvider), + }); + await waitForStorage(wallet); + + const chainStore = wallet.storage.getDefaultChainStore(); + const [storedHeader] = await mineHeadersChain('regtest', 1); + const transactions = new Map(); + chainStore.state.chainHeight = 42; + chainStore.state.lastSyncedBlockHeight = 40; + chainStore.state.blockHeaders = [storedHeader]; + chainStore.state.lastSyncedHeaderHeight = 41; + chainStore.state.headersMetadata = new Map([[storedHeader.hash, { height: 41 }]]); + chainStore.state.hashesByHeight = new Map([[41, storedHeader.hash]]); + chainStore.state.transactions = transactions; + const saveState = this.sinon.spy(wallet.storage, 'saveState'); + const storageKey = `wallet_${wallet.storage.currentWalletId}`; + + await wallet.createAccount({ + injectDefaultPlugins: false, + synchronize: false, + }); + + expect(saveState).to.have.been.calledOnce(); + expect(chainStore.state.blockHeaders).to.deep.equal([]); + expect(chainStore.state.lastSyncedHeaderHeight).to.equal(-1); + expect(chainStore.state.headersMetadata).to.deep.equal(new Map()); + expect(chainStore.state.hashesByHeight).to.deep.equal(new Map()); + expect(chainStore.state.chainHeight).to.equal(42); + expect(chainStore.state.lastSyncedBlockHeight).to.equal(40); + expect(chainStore.state.transactions).to.equal(transactions); + expect(blockHeadersProvider.spvChain.hashesByHeight.has(0)).to.equal(true); + + const normalizedStorage = await wallet.storage.adapter.getItem(storageKey); + const normalizedChain = normalizedStorage.chains[wallet.storage.currentNetwork]; + expect(normalizedChain.blockHeaders).to.deep.equal([]); + expect(normalizedChain.lastSyncedHeaderHeight).to.equal(-1); + + const worker = new BlockHeadersSyncWorker({ executeOnStart: false }); + worker.logger = logger; + worker.storage = wallet.storage; + worker.transport = wallet.transport; + worker.parentEvents = new EventEmitter(); + this.sinon.stub(worker, 'scheduleProgressUpdate'); + + expect(worker.getStartBlockHeight()).to.equal(1); + + const genesisChain = await mineHeadersChain('regtest', 2); + blockHeadersProvider.spvChain.addHeaders(genesisChain); + worker.historicalChainUpdateHandler(); + await wallet.storage.saveState(); + + expect(chainStore.state.lastSyncedHeaderHeight).to.equal(1); + expect(chainStore.state.hashesByHeight.has(0)).to.equal(true); + expect(chainStore.state.hashesByHeight.has(1)).to.equal(true); + + const rebuiltStorage = await wallet.storage.adapter.getItem(storageKey); + const rebuiltChain = rebuiltStorage.chains[wallet.storage.currentNetwork]; + expect(rebuiltChain.blockHeaders).to.have.lengthOf(2); + expect(rebuiltChain.lastSyncedHeaderHeight).to.equal(1); + }); + + it('should initialize and resume from the same minimal stored header context', async function shouldComposeStoredResumeContext() { + const blockHeadersProvider = new BlockHeadersProvider({ network: 'regtest' }); + wallet = new Wallet({ + mnemonic: null, + network: 'regtest', + transport: new TestTransport(blockHeadersProvider), + }); + await waitForStorage(wallet); + + const storedHeaders = await mineHeadersChain('regtest', 2); + const chainStore = wallet.storage.getDefaultChainStore(); + chainStore.state.blockHeaders = storedHeaders; + chainStore.state.lastSyncedHeaderHeight = 42; + const saveState = this.sinon.stub(wallet.storage, 'saveState').resolves(true); + + await wallet.createAccount({ + injectDefaultPlugins: false, + synchronize: false, + }); + + const worker = new BlockHeadersSyncWorker({ executeOnStart: false }); + worker.logger = logger; + worker.storage = wallet.storage; + + expect(saveState).to.have.not.been.called; + expect(blockHeadersProvider.spvChain.hashesByHeight.has(41)).to.equal(true); + expect(worker.getStartBlockHeight()).to.equal(42); + }); + + it('should initialize genesis when stored headers imply a negative first height', async function shouldRejectNegativeFirstHeight() { + const blockHeadersProvider = { + initializeChainWith: this.sinon.stub().resolves(), + }; + wallet = new Wallet({ + mnemonic: null, + network: 'regtest', + transport: new TestTransport(blockHeadersProvider), + }); + await waitForStorage(wallet); + + const chainStore = wallet.storage.getDefaultChainStore(); + chainStore.state.blockHeaders = [{ hash: 'header-0' }, { hash: 'header-1' }]; + chainStore.state.lastSyncedHeaderHeight = 0; + + await wallet.createAccount({ + injectDefaultPlugins: false, + synchronize: false, + }); + + expect(blockHeadersProvider.initializeChainWith) + .to.have.been.calledOnceWithExactly([], -1); + }); + + it('should reject malformed selected stored headers', async () => { + const blockHeadersProvider = new BlockHeadersProvider({ network: 'regtest' }); + wallet = new Wallet({ + mnemonic: null, + network: 'regtest', + transport: new TestTransport(blockHeadersProvider), + }); + await waitForStorage(wallet); + + const chainStore = wallet.storage.getDefaultChainStore(); + chainStore.state.blockHeaders = [{ malformed: true }, { malformed: true }]; + chainStore.state.lastSyncedHeaderHeight = 42; + + await expect(wallet.createAccount({ + injectDefaultPlugins: false, + synchronize: false, + })).to.be.rejected(); + + expect(wallet.accounts).to.deep.equal([]); + }); + + it('should serialize normalization, persistence, and provider initialization', async function shouldSerializeInitialization() { + let resolveSave; + const deferredSave = new Promise((resolve) => { + resolveSave = resolve; + }); + const blockHeadersProvider = { + initializeChainWith: this.sinon.stub().resolves(), + }; + wallet = new Wallet({ + mnemonic: null, + network: 'regtest', + transport: new TestTransport(blockHeadersProvider), + }); + await waitForStorage(wallet); + + const chainStore = wallet.storage.getDefaultChainStore(); + chainStore.state.blockHeaders = [{ hash: 'stored-header-42' }]; + chainStore.state.lastSyncedHeaderHeight = 42; + chainStore.state.headersMetadata.set('stored-header-42', { height: 42 }); + chainStore.state.hashesByHeight.set(42, 'stored-header-42'); + const saveState = this.sinon.stub(wallet.storage, 'saveState').returns(deferredSave); + + const accountPromises = [ + wallet.createAccount({ injectDefaultPlugins: false, synchronize: false }), + wallet.createAccount({ injectDefaultPlugins: false, synchronize: false }), + ]; + + await new Promise((resolve) => { + setTimeout(resolve, 0); + }); + + try { + expect(chainStore.state.lastSyncedHeaderHeight).to.equal(-1); + expect(saveState).to.have.been.calledOnce(); + expect(blockHeadersProvider.initializeChainWith).to.have.not.been.called; + } finally { + resolveSave(); + } + + await Promise.all(accountPromises); + + expect(blockHeadersProvider.initializeChainWith) + .to.have.been.calledOnceWithExactly([], -1); + }); + + it('should reject all concurrent accounts when normalized state cannot be saved', async function shouldRejectFailedNormalizationSave() { + const saveError = new Error('Header state save failed'); + const blockHeadersProvider = { + initializeChainWith: this.sinon.stub().resolves(), + }; + wallet = new Wallet({ + mnemonic: null, + network: 'regtest', + transport: new TestTransport(blockHeadersProvider), + }); + await waitForStorage(wallet); + + const chainStore = wallet.storage.getDefaultChainStore(); + chainStore.state.blockHeaders = [{ hash: 'stored-header-42' }]; + chainStore.state.lastSyncedHeaderHeight = 42; + const saveState = this.sinon.stub(wallet.storage, 'saveState').rejects(saveError); + + const accountPromises = [ + wallet.createAccount({ injectDefaultPlugins: false, synchronize: false }), + wallet.createAccount({ injectDefaultPlugins: false, synchronize: false }), + ]; + + try { + const results = await Promise.allSettled(accountPromises); + + expect(results.map(({ status }) => status)) + .to.deep.equal(['rejected', 'rejected']); + results.forEach(({ reason }) => { + expect(reason).to.equal(saveError); + }); + + expect(blockHeadersProvider.initializeChainWith).to.have.not.been.called; + expect(wallet.accounts).to.deep.equal([]); + } finally { + saveState.restore(); + } + }); + + it('should normalize in memory without forcing persistence when autoSave is false', async function shouldPreserveDisabledPersistence() { + const blockHeadersProvider = { + initializeChainWith: this.sinon.stub().resolves(), + }; + wallet = new Wallet({ + mnemonic: null, + network: 'regtest', + storage: { + autoSave: false, + }, + transport: new TestTransport(blockHeadersProvider), + }); + await waitForStorage(wallet); + + const chainStore = wallet.storage.getDefaultChainStore(); + chainStore.state.blockHeaders = [{ hash: 'stored-header-42' }]; + chainStore.state.lastSyncedHeaderHeight = 42; + const setItem = this.sinon.spy(wallet.storage.adapter, 'setItem'); + + await wallet.createAccount({ + injectDefaultPlugins: false, + synchronize: false, + }); + + expect(chainStore.state.blockHeaders).to.deep.equal([]); + expect(chainStore.state.lastSyncedHeaderHeight).to.equal(-1); + expect(setItem).to.have.not.been.called; + expect(blockHeadersProvider.initializeChainWith) + .to.have.been.calledOnceWithExactly([], -1); + }); + + it('should share one provider initialization between concurrent accounts', async function shouldInitializeOnce() { + const blockHeadersProvider = { + initializeChainWith: this.sinon.stub().resolves(), + }; + wallet = new Wallet({ + mnemonic: null, + network: 'regtest', + transport: new TestTransport(blockHeadersProvider), + }); + + await Promise.all([ + wallet.createAccount({ + injectDefaultPlugins: false, + synchronize: false, + }), + wallet.createAccount({ + injectDefaultPlugins: false, + synchronize: false, + }), + ]); + + expect(blockHeadersProvider.initializeChainWith).to.have.been.calledOnce(); + }); + + it('should return one account for concurrent requests of the same index', async function shouldReturnOneAccount() { + let resolveInitialization; + const initializationPromise = new Promise((resolve) => { + resolveInitialization = resolve; + }); + const blockHeadersProvider = { + initializeChainWith: this.sinon.stub().returns(initializationPromise), + }; + wallet = new Wallet({ + mnemonic: null, + network: 'regtest', + transport: new TestTransport(blockHeadersProvider), + }); + + const accountPromises = [ + wallet.getAccount({ + index: 0, + injectDefaultPlugins: false, + synchronize: false, + }), + wallet.getAccount({ + index: 0, + injectDefaultPlugins: false, + synchronize: false, + }), + ]; + + await new Promise((resolve) => { + setTimeout(resolve, 0); + }); + resolveInitialization(); + + const [firstAccount, secondAccount] = await Promise.all(accountPromises); + expect(firstAccount).to.equal(secondAccount); + expect(wallet.accounts).to.deep.equal([firstAccount]); + }); + + it('should propagate provider initialization failure without creating an account', async function shouldRejectInitialization() { + const initializationError = new Error('SPV initialization failed'); + const blockHeadersProvider = { + initializeChainWith: this.sinon.stub().rejects(initializationError), + }; + wallet = new Wallet({ + mnemonic: null, + network: 'regtest', + transport: new TestTransport(blockHeadersProvider), + }); + + await expect(wallet.createAccount({ + injectDefaultPlugins: false, + synchronize: false, + })).to.be.rejectedWith(initializationError); + + expect(wallet.accounts).to.deep.equal([]); + }); + + it('should retry provider initialization after a failed attempt', async function shouldRetryInitialization() { + const initializationError = new Error('SPV initialization failed'); + const initializeChainWith = this.sinon.stub(); + initializeChainWith.onFirstCall().rejects(initializationError); + initializeChainWith.onSecondCall().resolves(); + wallet = new Wallet({ + mnemonic: null, + network: 'regtest', + transport: new TestTransport({ initializeChainWith }), + }); + + await expect(wallet.createAccount({ + injectDefaultPlugins: false, + synchronize: false, + })).to.be.rejectedWith(initializationError); + + const account = await wallet.createAccount({ + injectDefaultPlugins: false, + synchronize: false, + }); + + expect(account).to.exist(); + expect(initializeChainWith).to.have.been.calledTwice(); + }); + + it('should preserve online custom transports without a block headers provider', async () => { + wallet = new Wallet({ + mnemonic: null, + network: 'regtest', + transport: new TestTransportWithoutClient(), + }); + + const account = await wallet.createAccount({ + injectDefaultPlugins: false, + synchronize: false, + }); + + expect(account).to.exist(); + expect(wallet.accounts).to.deep.equal([account]); + }); + + it('should create an offline account without a transport', async () => { + wallet = new Wallet({ + mnemonic: null, + offlineMode: true, + }); + + const account = await wallet.createAccount({ + injectDefaultPlugins: false, + synchronize: false, + }); + + expect(account).to.exist(); + expect(wallet.transport).to.equal(undefined); + }); +}); diff --git a/packages/wallet-lib/src/types/Wallet/methods/getAccount.js b/packages/wallet-lib/src/types/Wallet/methods/getAccount.js index 8e3f33c7128..b82af599ecd 100644 --- a/packages/wallet-lib/src/types/Wallet/methods/getAccount.js +++ b/packages/wallet-lib/src/types/Wallet/methods/getAccount.js @@ -32,7 +32,21 @@ async function getAccount(accountOpts = defaultOpts) { const baseOpts = { index: accountIndex }; const opts = Object.assign(baseOpts, accountOpts); - return acc[0] || this.createAccount(opts); + if (acc[0]) { + return acc[0]; + } + + if (!this.accountCreationPromises) { + this.accountCreationPromises = new Map(); + } + + if (!this.accountCreationPromises.has(accountIndex)) { + const accountCreationPromise = Promise.resolve(this.createAccount(opts)) + .finally(() => this.accountCreationPromises.delete(accountIndex)); + this.accountCreationPromises.set(accountIndex, accountCreationPromise); + } + + return this.accountCreationPromises.get(accountIndex); } module.exports = getAccount; diff --git a/packages/wasm-sdk/Cargo.toml b/packages/wasm-sdk/Cargo.toml index 47b08485563..179eea2bf4d 100644 --- a/packages/wasm-sdk/Cargo.toml +++ b/packages/wasm-sdk/Cargo.toml @@ -93,7 +93,11 @@ dpp-json-convertible-derive = { path = "../rs-dpp-json-convertible-derive" } wasm-dpp2 = { path = "../wasm-dpp2" } [dev-dependencies] -wasm-sdk = { path = ".", features = ["mocks", "all-system-contracts"] } +# Self-dependency to enable `mocks` for the lib under test. Do not add +# system-contract features here: Cargo unifies them into the tested lib, so +# `all-system-contracts` would let the trusted context resolve every contract +# and silently defeat the tests that assert an uncompiled one is fetched. +wasm-sdk = { path = ".", features = ["mocks"] } tokio = { version = "1.40", features = ["macros", "rt-multi-thread"] } [package.metadata.wasm-pack] diff --git a/packages/wasm-sdk/package.json b/packages/wasm-sdk/package.json index b55f73e60c6..0d581f94571 100644 --- a/packages/wasm-sdk/package.json +++ b/packages/wasm-sdk/package.json @@ -32,7 +32,7 @@ "build:release": "./scripts/build-optimized.sh && node ./scripts/bundle.cjs", "test": "yarn run test:unit", "test:unit": "mocha --import ts-node/esm tests/unit/**/*.spec.ts && karma start ./tests/karma/karma.conf.cjs --single-run", - "test:functional": "mocha --import ts-node/esm --require tests/functional-bootstrap.cjs tests/functional/**/*.spec.ts && karma start ./tests/karma/karma.functional.conf.cjs --single-run", + "test:functional": "mocha --import ts-node/esm --require tests/functional-bootstrap.cjs --require tests/functional-readiness.mjs tests/functional/**/*.spec.ts && karma start ./tests/karma/karma.functional.conf.cjs --single-run", "lint": "eslint tests/**/*.ts" }, "ultra": { diff --git a/packages/wasm-sdk/src/context_provider.rs b/packages/wasm-sdk/src/context_provider.rs index fce91103031..cd17e74aa3c 100644 --- a/packages/wasm-sdk/src/context_provider.rs +++ b/packages/wasm-sdk/src/context_provider.rs @@ -220,6 +220,14 @@ impl WasmTrustedContext { } impl WasmTrustedContext { + /// Refresh quorum keys before proof verification. + pub(crate) async fn refresh_quorums(&self) -> Result<(), WasmSdkError> { + self.inner + .refresh_quorum_caches() + .await + .map_err(|e| WasmSdkError::generic(format!("Failed to refresh quorums: {}", e))) + } + /// Shared constructor used by every `prefetch*` factory. When `base_url` /// is `Some`, it overrides the default URL derived from `network` + /// `devnet_name` (the validator inside `new_with_url` still runs). @@ -307,20 +315,29 @@ impl WasmTrustedContext { /// Build a `WasmTrustedContext` with the given `discovered_addresses` for /// use in unit tests. The inner provider is constructed against a local - /// loopback URL — its quorum/contract caches are unused by the tests and - /// the URL is never dialled — so we get a real `Arc` - /// without any network side effects. + /// loopback URL that these tests do not dial, providing a real + /// `Arc` without network side effects. #[cfg(test)] pub(crate) fn for_testing( discovered_addresses: Vec, + ) -> WasmTrustedContext { + Self::for_testing_with_url(discovered_addresses, "http://127.0.0.1:22444".to_string()) + } + + /// Build a test context whose provider reads from a controllable endpoint. + #[cfg(test)] + pub(crate) fn for_testing_with_url( + discovered_addresses: Vec, + base_url: String, ) -> WasmTrustedContext { let cache_size = std::num::NonZeroUsize::new(1).unwrap(); let inner = rs_sdk_trusted_context_provider::TrustedHttpContextProvider::new_with_url( dash_sdk::dpp::dashcore::Network::Regtest, - "http://127.0.0.1:22444".to_string(), + base_url, cache_size, ) - .expect("loopback URL must construct a TrustedHttpContextProvider"); + .expect("test URL must construct a TrustedHttpContextProvider") + .with_refetch_if_not_found(false); WasmTrustedContext { inner: Arc::new(inner), diff --git a/packages/wasm-sdk/src/sdk.rs b/packages/wasm-sdk/src/sdk.rs index 83c87428f8c..312531d30a0 100644 --- a/packages/wasm-sdk/src/sdk.rs +++ b/packages/wasm-sdk/src/sdk.rs @@ -98,17 +98,13 @@ impl WasmSdk { } } - /// Fetch a contract, checking cache first - pub(crate) async fn get_or_fetch_contract( + /// Fetch a proved contract and replace its trusted-context cache entry. + pub(crate) async fn refresh_contract( &self, contract_id: dash_sdk::platform::Identifier, ) -> Result { use dash_sdk::platform::Fetch; - if let Some(cached) = self.get_cached_contract(&contract_id) { - return Ok((*cached).clone()); - } - let contract = dash_sdk::platform::DataContract::fetch(self.as_ref(), contract_id) .await? .ok_or_else(|| crate::error::WasmSdkError::not_found("Data contract not found"))?; @@ -118,6 +114,18 @@ impl WasmSdk { Ok(contract) } + /// Fetch a contract, checking cache first + pub(crate) async fn get_or_fetch_contract( + &self, + contract_id: dash_sdk::platform::Identifier, + ) -> Result { + if let Some(cached) = self.get_cached_contract(&contract_id) { + return Ok((*cached).clone()); + } + + self.refresh_contract(contract_id).await + } + /// Remove a contract from the cache pub(crate) fn remove_cached_contract( &self, @@ -465,9 +473,27 @@ impl WasmSdkBuilder { } } +#[cfg(test)] +impl WasmSdk { + /// Pair a mock `Sdk` with a trusted context so tests in sibling modules can + /// drive the paths that read the contract and quorum caches. + pub(crate) fn new_for_testing(sdk: Sdk, trusted_context: Option) -> Self { + Self { + sdk, + trusted_context, + } + } +} + #[cfg(test)] mod tests { use super::*; + use dash_sdk::dpp::data_contract::accessors::v0::{ + DataContractV0Getters, DataContractV0Setters, + }; + use dash_sdk::dpp::prelude::DataContract; + use dash_sdk::dpp::system_data_contracts::{load_system_data_contract, SystemDataContract}; + use dash_sdk::dpp::version::PlatformVersion; use dash_sdk::sdk::Uri; use rs_dapi_client::Address; @@ -496,6 +522,22 @@ mod tests { .collect() } + /// Build the fixture at the SDK's own platform version. A contract fetched + /// through the mock is deserialized at that version, and document types + /// differ between versions, so a fixture pinned to `latest` would never + /// compare equal to the round-tripped one. + fn custom_contract( + id_byte: u8, + version: u32, + platform_version: &PlatformVersion, + ) -> DataContract { + let mut contract = load_system_data_contract(SystemDataContract::DPNS, platform_version) + .expect("DPNS contract fixture should load"); + contract.set_id(dash_sdk::platform::Identifier::new([id_byte; 32])); + contract.set_version(version); + contract + } + #[test] fn with_addresses_sets_user_addresses_flag() { let b = user_builder(); @@ -619,4 +661,99 @@ mod tests { let b = user_builder().with_trusted_context(&ctx); assert!(b.has_user_addresses()); } + + #[tokio::test] + async fn refresh_contract_caches_a_missing_custom_contract() { + let mut inner_sdk = Sdk::new_mock(); + let expected = custom_contract(0x44, 1, inner_sdk.version()); + let contract_id = expected.id(); + inner_sdk + .mock() + .expect_fetch(contract_id, Some(expected.clone())) + .await + .expect("mock contract response should be configured"); + + let sdk = WasmSdk { + sdk: inner_sdk, + trusted_context: Some(WasmTrustedContext::for_testing(vec![])), + }; + assert!(sdk.get_cached_contract(&contract_id).is_none()); + + let refreshed = sdk + .refresh_contract(contract_id) + .await + .expect("proved contract fetch should succeed"); + + assert_eq!(refreshed, expected); + assert_eq!( + sdk.get_cached_contract(&contract_id) + .expect("refreshed contract should be cached") + .as_ref(), + &expected, + ); + } + + #[tokio::test] + async fn refresh_contract_replaces_a_stale_cached_version() { + let mut inner_sdk = Sdk::new_mock(); + let stale = custom_contract(0x55, 1, inner_sdk.version()); + let expected = custom_contract(0x55, 2, inner_sdk.version()); + let contract_id = expected.id(); + let context = WasmTrustedContext::for_testing(vec![]); + context.add_known_contract(stale); + + inner_sdk + .mock() + .expect_fetch(contract_id, Some(expected.clone())) + .await + .expect("mock contract response should be configured"); + + let sdk = WasmSdk { + sdk: inner_sdk, + trusted_context: Some(context), + }; + let refreshed = sdk + .refresh_contract(contract_id) + .await + .expect("proved contract refresh should succeed"); + + assert_eq!(refreshed.version(), 2); + assert_eq!( + sdk.get_cached_contract(&contract_id) + .expect("latest contract should replace stale cache") + .version(), + 2, + ); + } + + #[tokio::test] + async fn refresh_contract_propagates_proved_absence() { + let contract_id = dash_sdk::platform::Identifier::new([0x66; 32]); + let mut inner_sdk = Sdk::new_mock(); + inner_sdk + .mock() + .expect_fetch(contract_id, None as Option) + .await + .expect("mock absence response should be configured"); + + let sdk = WasmSdk { + sdk: inner_sdk, + trusted_context: Some(WasmTrustedContext::for_testing(vec![])), + }; + + assert!(sdk.refresh_contract(contract_id).await.is_err()); + assert!(sdk.get_cached_contract(&contract_id).is_none()); + } + + #[tokio::test] + async fn refresh_contract_propagates_fetch_errors() { + let contract_id = dash_sdk::platform::Identifier::new([0x77; 32]); + let sdk = WasmSdk { + sdk: Sdk::new_mock(), + trusted_context: Some(WasmTrustedContext::for_testing(vec![])), + }; + + assert!(sdk.refresh_contract(contract_id).await.is_err()); + assert!(sdk.get_cached_contract(&contract_id).is_none()); + } } diff --git a/packages/wasm-sdk/src/state_transitions/broadcast.rs b/packages/wasm-sdk/src/state_transitions/broadcast.rs index 60f779448c2..b2b2ca30d9e 100644 --- a/packages/wasm-sdk/src/state_transitions/broadcast.rs +++ b/packages/wasm-sdk/src/state_transitions/broadcast.rs @@ -6,15 +6,83 @@ use crate::error::WasmSdkError; use crate::sdk::WasmSdk; use crate::settings::{parse_put_settings, PutSettingsJs}; +use dash_sdk::dpp::platform_value::Identifier; +use dash_sdk::dpp::state_transition::batch_transition::accessors::DocumentsBatchTransitionAccessorsV0; use dash_sdk::dpp::state_transition::proof_result::StateTransitionProofResult; use dash_sdk::dpp::state_transition::StateTransition; use dash_sdk::platform::transition::broadcast::BroadcastStateTransition; +use dash_sdk::platform::ContextProvider; +use std::collections::BTreeSet; use wasm_bindgen::prelude::*; use wasm_dpp2::state_transitions::proof_result::{ convert_proof_result, StateTransitionProofResultTypeJs, }; use wasm_dpp2::StateTransitionWasm; +fn referenced_contract_ids(state_transition: &StateTransition) -> BTreeSet { + match state_transition { + StateTransition::Batch(batch_transition) => batch_transition + .transitions_iter() + .map(|transition| transition.data_contract_id()) + .collect(), + _ => BTreeSet::new(), + } +} + +impl WasmSdk { + /// Whether the context provider can already supply this contract, either + /// from its cache or from a definition compiled into the SDK. + fn can_resolve_contract(&self, contract_id: Identifier) -> bool { + self.trusted_context() + .and_then(|context| { + context + // `WasmSdk::version` is the exported protocol number; the + // platform version comes from the inner SDK. + .get_data_contract(&contract_id, self.as_ref().version()) + .ok() + .flatten() + }) + .is_some() + } +} + +impl WasmSdk { + async fn prepare_state_transition_context( + &self, + state_transition: &StateTransition, + ) -> Result<(), WasmSdkError> { + if let Some(context) = self.trusted_context() { + if let Err(error) = context.refresh_quorums().await { + tracing::warn!( + error = %error, + "Failed to refresh trusted quorum cache before proof verification; using cached keys" + ); + } + } + + for contract_id in referenced_contract_ids(state_transition) { + // A contract the provider already resolves is left alone: fetching + // it would let a node-supplied copy shadow a cached or compiled-in + // definition, which decides verification when proofs are disabled. + // Which system contracts are compiled in depends on cargo features + // — the withdrawals contract is not among the wasm defaults — so + // this asks the provider rather than assuming. + if self.can_resolve_contract(contract_id) { + continue; + } + + // Only an unresolvable contract reaches the network. This runs + // after the transition is broadcast, where refetching what we + // already hold would let one transient failure discard a result + // the network already accepted. A contract that can be neither + // resolved nor fetched stays fatal: the proof needs it. + self.refresh_contract(contract_id).await?; + } + + Ok(()) + } +} + #[wasm_bindgen] impl WasmSdk { /// Broadcasts a state transition to the network. @@ -61,23 +129,25 @@ impl WasmSdk { ) -> Result { let st: StateTransition = state_transition.into(); let put_settings = parse_put_settings(settings)?; + self.prepare_state_transition_context(&st).await?; + // Preserve the error kind: callers distinguish `ExecutionNotProved`, + // which reports that this transition family has no execution proof + // rather than that anything went wrong, from a genuine failure. let result = st .wait_for_response::(self.as_ref(), put_settings) .await - .map_err(|e| { - WasmSdkError::generic(format!("Failed to wait for state transition result: {}", e)) - })?; + .map_err(WasmSdkError::from)?; convert_proof_result(result).map_err(WasmSdkError::from) } /// Broadcasts a state transition and waits for the result. /// - /// This method broadcasts the transition and waits for confirmation from the network. - /// Returns once the transition has been processed or fails. - /// This is equivalent to calling `broadcastStateTransition` followed by - /// `waitForResponse`. + /// This method prepares proof context, broadcasts the transition, and waits + /// for confirmation from the network. Returns once the transition has been + /// processed or fails. Unlike separate broadcast and wait calls, proof + /// context preparation happens before broadcasting. /// /// @param stateTransition - The state transition to broadcast /// @param settings - Optional put settings (retries, timeout, waitTimeoutMs) @@ -90,6 +160,7 @@ impl WasmSdk { ) -> Result { let st: StateTransition = state_transition.into(); let put_settings = parse_put_settings(settings)?; + self.prepare_state_transition_context(&st).await?; let result = st .broadcast_and_wait::(self.as_ref(), put_settings) @@ -121,6 +192,7 @@ impl WasmSdk { ) -> Result { let st: StateTransition = state_transition.into(); let put_settings = parse_put_settings(settings)?; + self.prepare_state_transition_context(&st).await?; let result = st .wait_for_affected_state::(self.as_ref(), put_settings) @@ -147,6 +219,7 @@ impl WasmSdk { ) -> Result { let st: StateTransition = state_transition.into(); let put_settings = parse_put_settings(settings)?; + self.prepare_state_transition_context(&st).await?; let result = st .broadcast_and_wait_for_affected_state::( @@ -159,3 +232,286 @@ impl WasmSdk { convert_proof_result(result).map_err(WasmSdkError::from) } } + +#[cfg(test)] +mod tests { + use super::*; + use crate::context_provider::WasmTrustedContext; + use crate::sdk::WasmSdkBuilder; + use dash_sdk::dpp::state_transition::batch_transition::batched_transition::document_base_transition::v0::DocumentBaseTransitionV0; + use dash_sdk::dpp::state_transition::batch_transition::batched_transition::document_base_transition::DocumentBaseTransition; + use dash_sdk::dpp::state_transition::batch_transition::batched_transition::document_delete_transition::v0::DocumentDeleteTransitionV0; + use dash_sdk::dpp::state_transition::batch_transition::batched_transition::document_delete_transition::DocumentDeleteTransition; + use dash_sdk::dpp::state_transition::batch_transition::batched_transition::document_transition::DocumentTransition; + use dash_sdk::dpp::state_transition::batch_transition::{BatchTransition, BatchTransitionV0}; + use dash_sdk::dpp::state_transition::identity_topup_transition::v0::IdentityTopUpTransitionV0; + use dash_sdk::dpp::state_transition::identity_topup_transition::IdentityTopUpTransition; + use dash_sdk::dpp::data_contract::accessors::v0::{ + DataContractV0Getters, DataContractV0Setters, + }; + use dash_sdk::dpp::platform_value::BinaryData; + use dash_sdk::dpp::prelude::DataContract; + use dash_sdk::dpp::system_data_contracts::{load_system_data_contract, SystemDataContract}; + use dash_sdk::dpp::version::PlatformVersion; + use dash_sdk::Sdk; + use std::io::{BufRead, BufReader, Write}; + use std::net::{TcpListener, TcpStream}; + use std::thread; + use std::time::{Duration, Instant}; + + fn accept_before(listener: &TcpListener, deadline: Instant) -> TcpStream { + loop { + match listener.accept() { + Ok((stream, _)) => return stream, + Err(error) + if error.kind() == std::io::ErrorKind::WouldBlock + && Instant::now() < deadline => + { + thread::sleep(Duration::from_millis(10)); + } + Err(error) => panic!("accept quorum request: {}", error), + } + } + } + + fn spawn_rotated_quorum_endpoint() -> (String, thread::JoinHandle<()>) { + let listener = TcpListener::bind("127.0.0.1:0").expect("bind mock quorum endpoint"); + listener + .set_nonblocking(true) + .expect("make mock endpoint bounded"); + let address = listener.local_addr().expect("read mock endpoint address"); + + let handle = thread::spawn(move || { + let current = serde_json::json!({ + "success": true, + "data": [{ + "quorum_hash": hex::encode([0x88; 32]), + "key": hex::encode([0x98; 48]), + "height": 1, + "valid_members_count": 3 + }] + }) + .to_string(); + let previous = serde_json::json!({ + "success": true, + "data": { + "height": 1, + "quorums": [] + } + }) + .to_string(); + + for (expected_path, body) in [("/quorums", current), ("/previous", previous)] { + let mut stream = accept_before(&listener, Instant::now() + Duration::from_secs(5)); + let mut reader = + BufReader::new(stream.try_clone().expect("clone quorum request stream")); + let mut request_line = String::new(); + reader + .read_line(&mut request_line) + .expect("read quorum request line"); + assert_eq!(request_line.split_whitespace().nth(1), Some(expected_path)); + + loop { + let mut header = String::new(); + reader + .read_line(&mut header) + .expect("read quorum request header"); + if header == "\r\n" || header.is_empty() { + break; + } + } + + write!( + stream, + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", + body.len(), + body + ) + .expect("write quorum response"); + stream.flush().expect("flush quorum response"); + } + }); + + (format!("http://{}", address), handle) + } + + fn delete_transition(contract_id: Identifier, nonce: u64) -> DocumentTransition { + DocumentTransition::Delete(DocumentDeleteTransition::V0(DocumentDeleteTransitionV0 { + base: DocumentBaseTransition::V0(DocumentBaseTransitionV0 { + id: Identifier::new([nonce as u8; 32]), + identity_contract_nonce: nonce, + document_type_name: "note".to_string(), + data_contract_id: contract_id, + }), + })) + } + + #[test] + fn should_collect_each_referenced_contract_of_a_document_batch_once_in_order() { + let first_contract_id = Identifier::new([0x11; 32]); + let second_contract_id = Identifier::new([0x22; 32]); + let state_transition = StateTransition::Batch(BatchTransition::V0(BatchTransitionV0 { + owner_id: Identifier::new([0x33; 32]), + transitions: vec![ + delete_transition(first_contract_id, 1), + delete_transition(second_contract_id, 2), + delete_transition(first_contract_id, 3), + ], + user_fee_increase: 0, + signature_public_key_id: 0, + signature: BinaryData::default(), + })); + + assert_eq!( + referenced_contract_ids(&state_transition) + .into_iter() + .collect::>(), + vec![first_contract_id, second_contract_id], + ); + } + + #[test] + fn should_not_prepare_contracts_for_non_batch_transition() { + let state_transition = StateTransition::IdentityTopUp(IdentityTopUpTransition::V0( + IdentityTopUpTransitionV0::default(), + )); + + assert!(referenced_contract_ids(&state_transition).is_empty()); + } + + #[tokio::test] + async fn should_refresh_rotated_quorums_through_proof_preparation() { + let (base_url, server) = spawn_rotated_quorum_endpoint(); + let context = WasmTrustedContext::for_testing_with_url(vec![], base_url); + let sdk = WasmSdkBuilder::new_local() + .with_trusted_context(&context) + .build() + .expect("build local SDK with trusted context"); + let state_transition = StateTransition::IdentityTopUp(IdentityTopUpTransition::V0( + IdentityTopUpTransitionV0::default(), + )); + + assert!(context.get_quorum_public_key(1, [0x88; 32], 1).is_err()); + sdk.prepare_state_transition_context(&state_transition) + .await + .expect("proof preparation must refresh quorums"); + assert_eq!( + context + .get_quorum_public_key(1, [0x88; 32], 1) + .expect("rotated quorum must be available after preparation"), + [0x98; 48] + ); + server.join().expect("mock quorum server must finish"); + } + + /// Build the fixture at the SDK's own platform version. A contract fetched + /// through the mock is deserialized at that version, and document types + /// differ between versions, so a fixture pinned to `latest` would never + /// compare equal to the round-tripped one. + fn custom_contract( + id_byte: u8, + version: u32, + platform_version: &PlatformVersion, + ) -> DataContract { + let mut contract = load_system_data_contract(SystemDataContract::DPNS, platform_version) + .expect("DPNS contract fixture should load"); + contract.set_id(Identifier::new([id_byte; 32])); + contract.set_version(version); + contract + } + + #[tokio::test] + async fn should_cache_a_referenced_contract_missing_from_the_context() { + let mut inner_sdk = Sdk::new_mock(); + let expected = custom_contract(0x66, 1, inner_sdk.version()); + let contract_id = expected.id(); + inner_sdk + .mock() + .expect_fetch(contract_id, Some(expected.clone())) + .await + .expect("mock contract response should be configured"); + + let context = WasmTrustedContext::for_testing(vec![]); + let sdk = WasmSdk::new_for_testing(inner_sdk, Some(context)); + let state_transition = StateTransition::Batch(BatchTransition::V0(BatchTransitionV0 { + owner_id: Identifier::new([0x33; 32]), + transitions: vec![delete_transition(contract_id, 1)], + user_fee_increase: 0, + signature_public_key_id: 0, + signature: BinaryData::default(), + })); + + assert!(sdk.get_cached_contract(&contract_id).is_none()); + sdk.prepare_state_transition_context(&state_transition) + .await + .expect("preparation must fetch the referenced contract"); + assert_eq!( + sdk.get_cached_contract(&contract_id) + .expect("referenced contract must be cached after preparation") + .as_ref(), + &expected, + ); + } + + #[tokio::test] + async fn should_fetch_a_referenced_system_contract_the_sdk_does_not_compile_in() { + // Which system contracts are compiled in is a build-time choice, and + // the withdrawals contract is not among the wasm defaults. The provider + // cannot serve it, so preparation has to fetch it like any other + // contract or document verification fails on an unknown contract. + let withdrawals_id = SystemDataContract::Withdrawals.id(); + let mut inner_sdk = Sdk::new_mock(); + let mut expected = load_system_data_contract(SystemDataContract::DPNS, inner_sdk.version()) + .expect("DPNS contract fixture should load"); + expected.set_id(withdrawals_id); + + inner_sdk + .mock() + .expect_fetch(withdrawals_id, Some(expected.clone())) + .await + .expect("mock contract response should be configured"); + + let sdk = + WasmSdk::new_for_testing(inner_sdk, Some(WasmTrustedContext::for_testing(vec![]))); + let state_transition = StateTransition::Batch(BatchTransition::V0(BatchTransitionV0 { + owner_id: Identifier::new([0x33; 32]), + transitions: vec![delete_transition(withdrawals_id, 1)], + user_fee_increase: 0, + signature_public_key_id: 0, + signature: BinaryData::default(), + })); + + sdk.prepare_state_transition_context(&state_transition) + .await + .expect("preparation must fetch a system contract the SDK cannot resolve"); + assert_eq!( + sdk.get_cached_contract(&withdrawals_id) + .expect("fetched system contract must be cached after preparation") + .as_ref(), + &expected, + ); + } + + #[tokio::test] + async fn should_keep_the_compiled_in_definition_of_a_referenced_system_contract() { + let dpns_id = SystemDataContract::DPNS.id(); + // The mock SDK has no response configured, so any fetch of this id fails + // the preparation and proves the system contract was not requested. + let sdk = WasmSdk::new_for_testing( + Sdk::new_mock(), + Some(WasmTrustedContext::for_testing(vec![])), + ); + let state_transition = StateTransition::Batch(BatchTransition::V0(BatchTransitionV0 { + owner_id: Identifier::new([0x33; 32]), + transitions: vec![delete_transition(dpns_id, 1)], + user_fee_increase: 0, + signature_public_key_id: 0, + signature: BinaryData::default(), + })); + + sdk.prepare_state_transition_context(&state_transition) + .await + .expect("preparation must skip compiled-in system contracts"); + assert!(sdk.get_cached_contract(&dpns_id).is_none()); + } +} diff --git a/packages/wasm-sdk/tests/functional-readiness.mjs b/packages/wasm-sdk/tests/functional-readiness.mjs new file mode 100644 index 00000000000..c994dec76ca --- /dev/null +++ b/packages/wasm-sdk/tests/functional-readiness.mjs @@ -0,0 +1,17 @@ +import init from '../dist/sdk.compressed.js'; +import { prefetchLocalReady } from './functional/helpers/trustedContext.ts'; + +const NETWORK_READY_TIMEOUT_MS = 600_000; +const HOOK_TIMEOUT_MS = 630_000; + +export const mochaHooks = { + async beforeAll() { + this.timeout(HOOK_TIMEOUT_MS); + + await init(); + const context = await prefetchLocalReady({ + timeoutMs: NETWORK_READY_TIMEOUT_MS, + }); + context.free(); + }, +}; diff --git a/packages/wasm-sdk/tests/functional/epochs-blocks.spec.ts b/packages/wasm-sdk/tests/functional/epochs-blocks.spec.ts index bd6191b1d65..4e5737e9151 100644 --- a/packages/wasm-sdk/tests/functional/epochs-blocks.spec.ts +++ b/packages/wasm-sdk/tests/functional/epochs-blocks.spec.ts @@ -7,6 +7,7 @@ describe('Epochs and Evonode Blocks', function describeEpochs() { let client: sdk.WasmSdk; let evonodeProTxHash: string; + let epochIndex: number; before(async () => { await init(); @@ -17,6 +18,10 @@ describe('Epochs and Evonode Blocks', function describeEpochs() { // Get the proTxHash from the node status const status = await client.getStatus(); evonodeProTxHash = status.node.proTxHash; + if (status.time.epoch === undefined) { + throw new Error('Platform status did not include the current epoch'); + } + epochIndex = Number(status.time.epoch); }); after(async () => { @@ -27,10 +32,7 @@ describe('Epochs and Evonode Blocks', function describeEpochs() { describe('getEpochsInfo()', () => { it('should get epochs info and finalized epochs', async () => { - // Get current epoch info - const current = await client.getCurrentEpoch().catch(() => null); - const currentIndex = current ? Number(current.index) : 0; - const start = Math.max(0, currentIndex - 5); + const start = Math.max(0, epochIndex - 5); const infos = await client.getEpochsInfo({ startEpoch: start, @@ -43,10 +45,7 @@ describe('Epochs and Evonode Blocks', function describeEpochs() { describe('getFinalizedEpochInfos()', () => { it('should get finalized epoch infos', async () => { - // Get current epoch info - const current = await client.getCurrentEpoch().catch(() => null); - const currentIndex = current ? Number(current.index) : 0; - const start = Math.max(0, currentIndex - 5); + const start = Math.max(0, epochIndex - 5); const finalized = await client.getFinalizedEpochInfos({ startEpoch: start, @@ -58,10 +57,6 @@ describe('Epochs and Evonode Blocks', function describeEpochs() { describe('getEvonodesProposedEpochBlocksByIds()', () => { it('should query evonode proposed blocks by ids', async () => { - // Get current epoch - const current = await client.getCurrentEpoch().catch(() => null); - const epochIndex = current ? Number(current.index) : 0; - // Query by specific IDs only if we have a proTxHash if (evonodeProTxHash) { const byIds = await client @@ -73,10 +68,6 @@ describe('Epochs and Evonode Blocks', function describeEpochs() { describe('getEvonodesProposedEpochBlocksByRange()', () => { it('should query evonode proposed blocks by range', async () => { - // Get current epoch - const current = await client.getCurrentEpoch().catch(() => null); - const epochIndex = current ? Number(current.index) : 0; - // Query by range (doesn't require a specific proTxHash) const byRange = await client.getEvonodesProposedEpochBlocksByRange({ epoch: epochIndex, @@ -88,9 +79,6 @@ describe('Epochs and Evonode Blocks', function describeEpochs() { describe('getEvonodesProposedEpochBlocksByIdsWithProofInfo()', () => { it('should query evonode proposed blocks by ids with proof', async () => { - const current = await client.getCurrentEpoch().catch(() => null); - const epochIndex = current ? Number(current.index) : 0; - // Get at least one proTxHash from the range query results const byRange = await client.getEvonodesProposedEpochBlocksByRange({ epoch: epochIndex, @@ -122,9 +110,6 @@ describe('Epochs and Evonode Blocks', function describeEpochs() { describe('getEvonodesProposedEpochBlocksByRangeWithProofInfo()', () => { it('should query evonode proposed blocks by range with proof', async () => { - const current = await client.getCurrentEpoch().catch(() => null); - const epochIndex = current ? Number(current.index) : 0; - const res = await client.getEvonodesProposedEpochBlocksByRangeWithProofInfo({ epoch: epochIndex, limit: 50, diff --git a/packages/wasm-sdk/tests/functional/helpers/trustedContext.ts b/packages/wasm-sdk/tests/functional/helpers/trustedContext.ts index 176ec2379fa..630f813b990 100644 --- a/packages/wasm-sdk/tests/functional/helpers/trustedContext.ts +++ b/packages/wasm-sdk/tests/functional/helpers/trustedContext.ts @@ -8,35 +8,56 @@ export interface PrefetchLocalReadyOptions { intervalMs?: number; } +export interface PrefetchLocalReadyDependencies { + prefetchLocal?: () => Promise; + now?: () => number; + sleep?: (durationMs: number) => Promise; +} + // Wait for the local dashmate network to be ready, then return a prefetched // WasmTrustedContext. Retries on any error from prefetchLocal() until the // timeout — masternodes can take time to reach status=ENABLED + versionCheck=success // after `yarn start`, and a single failed attempt would otherwise abort the suite. export async function prefetchLocalReady( options: PrefetchLocalReadyOptions = {}, + dependencies: PrefetchLocalReadyDependencies = {}, ): Promise { const timeoutMs = options.timeoutMs ?? DEFAULT_TIMEOUT_MS; const intervalMs = options.intervalMs ?? DEFAULT_INTERVAL_MS; - const deadline = Date.now() + timeoutMs; + const prefetchLocal = dependencies.prefetchLocal + ?? (() => sdk.WasmTrustedContext.prefetchLocal()); + const now = dependencies.now ?? (() => performance.now()); + const sleep = dependencies.sleep ?? ((durationMs: number) => ( + new Promise((resolve) => { + setTimeout(resolve, durationMs); + }) + )); + const deadline = now() + timeoutMs; let lastError: unknown; let attempts = 0; - while (Date.now() < deadline) { + while (now() < deadline) { attempts += 1; try { - return await sdk.WasmTrustedContext.prefetchLocal(); + return await prefetchLocal(); } catch (error) { lastError = error; - const remaining = deadline - Date.now(); - if (remaining > intervalMs) { - await new Promise((resolve) => { - setTimeout(resolve, intervalMs); - }); + const remaining = deadline - now(); + if (remaining > 0) { + await sleep(Math.min(intervalMs, remaining)); } } } - const message = lastError instanceof Error ? lastError.message : String(lastError); + let message: string; + if (lastError instanceof Error) { + message = lastError.message; + } else if (lastError && typeof lastError === 'object' && 'message' in lastError) { + message = String(lastError.message); + } else { + message = String(lastError); + } + throw new Error( `prefetchLocalReady: local network not ready after ${timeoutMs}ms (${attempts} attempts); last error: ${message}`, ); diff --git a/packages/wasm-sdk/tests/functional/transitions/documents.spec.ts b/packages/wasm-sdk/tests/functional/transitions/documents.spec.ts index 63843527d05..762560aa1df 100644 --- a/packages/wasm-sdk/tests/functional/transitions/documents.spec.ts +++ b/packages/wasm-sdk/tests/functional/transitions/documents.spec.ts @@ -706,14 +706,29 @@ describe('Document State Transitions', function describeDocumentStateTransitions tokenPaymentInfo: makeTokenPaymentInfo(3n), }); - const purchasedDocument = await client.getDocument( + const purchaseDeadline = Date.now() + 30000; + let purchasedDocument = await client.getDocument( tokenPaidContractId, 'tokenPaidListing', tokenPaidDocumentId, ); + while ( + purchasedDocument?.ownerId.toString() !== testData.identityId3 + && Date.now() < purchaseDeadline + ) { + await new Promise((resolve) => { setTimeout(resolve, 500); }); + purchasedDocument = await client.getDocument( + tokenPaidContractId, + 'tokenPaidListing', + tokenPaidDocumentId, + ); + } expect(purchasedDocument).to.exist(); - expect(purchasedDocument.ownerId.toString()).to.equal(testData.identityId3); + expect(purchasedDocument.ownerId.toString()).to.equal( + testData.identityId3, + `owner did not converge in 30s; last owner: ${purchasedDocument.ownerId}`, + ); await expectTokenBalance(testData.identityId2, tokenPaidTokenId, 26n); await expectTokenBalance(testData.identityId3, tokenPaidTokenId, 47n); await expectTokenBalance(testData.identityId, tokenPaidTokenId, 927n); diff --git a/packages/wasm-sdk/tests/unit/trusted-context-ready.spec.ts b/packages/wasm-sdk/tests/unit/trusted-context-ready.spec.ts new file mode 100644 index 00000000000..e525d74b27c --- /dev/null +++ b/packages/wasm-sdk/tests/unit/trusted-context-ready.spec.ts @@ -0,0 +1,42 @@ +import { expect } from './helpers/chai.ts'; +import { prefetchLocalReady } from '../functional/helpers/trustedContext.ts'; + +describe('prefetchLocalReady()', () => { + it('should sleep through the deadline and report an object rejection message', async () => { + const sleepDurations: number[] = []; + const timestamps = [0, 0, 9, 10]; + let attempts = 0; + let thrownError: unknown; + + try { + await prefetchLocalReady( + { + timeoutMs: 10, + intervalMs: 4, + }, + { + prefetchLocal: () => { + attempts += 1; + // wasm-bindgen rejects with structured objects that are not Error instances. + // eslint-disable-next-line prefer-promise-reject-errors + return Promise.reject({ message: 'quorum is not ready' }); + }, + now: () => timestamps.shift() ?? 10, + sleep: async (durationMs) => { + sleepDurations.push(durationMs); + }, + }, + ); + } catch (error) { + thrownError = error; + } + + expect(attempts).to.equal(1); + expect(sleepDurations).to.deep.equal([1]); + expect(thrownError).to.be.instanceOf(Error); + expect((thrownError as Error).message).to.equal( + 'prefetchLocalReady: local network not ready after 10ms ' + + '(1 attempts); last error: quorum is not ready', + ); + }); +});