diff --git a/scenarios/create-dataset-smoke.mjs b/scenarios/create-dataset-smoke.mjs new file mode 100644 index 0000000..fa91aac --- /dev/null +++ b/scenarios/create-dataset-smoke.mjs @@ -0,0 +1,201 @@ +// Copied into synapse-sdk/utils before execution; imports are relative to that destination. +import { readFileSync } from 'fs' +import { homedir } from 'os' +import { join } from 'path' +import { http as viemHttp, maxUint256 } from 'viem' +import { privateKeyToAccount } from 'viem/accounts' +import { Synapse } from '../packages/synapse-sdk/src/index.ts' +import * as ERC20 from '../packages/synapse-core/src/erc20/index.ts' +import * as Pay from '../packages/synapse-core/src/pay/index.ts' +import * as SP from '../packages/synapse-core/src/sp/index.ts' +import { toChain, validateDevnetInfo } from '../packages/synapse-core/src/devnet/index.ts' + +const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms)) + +function assert(condition, message) { + if (!condition) { + throw new Error(message) + } +} + +async function waitForDataSet(synapse, dataSetId) { + let lastCount = 0 + for (let attempt = 1; attempt <= 15; attempt++) { + const dataSets = await synapse.storage.findDataSets() + lastCount = dataSets.length + const dataSet = dataSets.find((item) => item.dataSetId === dataSetId || item.pdpVerifierDataSetId === dataSetId) + if (dataSet != null) { + return dataSet + } + await sleep(2000) + } + throw new Error(`Created data set ${dataSetId} was not returned by findDataSets; last count=${lastCount}`) +} + +async function waitForTransactionReceipt(client, hash, label) { + let lastError = null + for (let attempt = 1; attempt <= 90; attempt++) { + try { + const receipt = await client.request({ + method: 'eth_getTransactionReceipt', + params: [hash], + }) + if (receipt != null) { + return receipt + } + } catch (error) { + lastError = error + } + + await sleep(1000) + } + + const suffix = lastError == null ? '' : `; last error: ${lastError.message}` + throw new Error(`${label} transaction ${hash} was not confirmed${suffix}`) +} + +async function prepareWithPlainErc20(synapse, context) { + const { costs, transaction } = await synapse.storage.prepare({ + context, + dataSize: 1n, + }) + console.log(`Prepared account: ready=${costs.ready}, depositNeeded=${costs.depositNeeded}`) + if (transaction == null) { + return + } + + if (costs.depositNeeded > 0n) { + const approveHash = await ERC20.approve(synapse.client, { + amount: costs.depositNeeded, + }) + console.log(`ERC20 approve tx submitted: ${approveHash}`) + const approved = await waitForTransactionReceipt(synapse.client, approveHash, 'ERC20 approve') + console.log(`ERC20 approve confirmed in block ${approved.blockNumber}`) + + const depositHash = await Pay.deposit(synapse.client, { + amount: costs.depositNeeded, + }) + console.log(`FilecoinPay deposit tx submitted: ${depositHash}`) + const deposited = await waitForTransactionReceipt(synapse.client, depositHash, 'FilecoinPay deposit') + console.log(`FilecoinPay deposit confirmed in block ${deposited.blockNumber}`) + } + + if (costs.needsFwssMaxApproval) { + const operatorApprovalHash = await Pay.setOperatorApproval(synapse.client, { + approve: true, + rateAllowance: maxUint256, + lockupAllowance: maxUint256, + }) + console.log(`FWSS operator approval tx submitted: ${operatorApprovalHash}`) + const approvedOperator = await waitForTransactionReceipt( + synapse.client, + operatorApprovalHash, + 'FWSS operator approval' + ) + console.log(`FWSS operator approval confirmed in block ${approvedOperator.blockNumber}`) + } +} + +async function main() { + const devnetInfoPath = + process.env.DEVNET_INFO_PATH || join(homedir(), '.foc-devnet', 'state', 'latest', 'devnet-info.json') + const userIndex = Number(process.env.DEVNET_USER_INDEX || '1') + const raw = JSON.parse(readFileSync(devnetInfoPath, 'utf8')) + const devnetInfo = validateDevnetInfo(raw) + const { info } = devnetInfo + + assert(Number.isInteger(userIndex), `DEVNET_USER_INDEX must be an integer; got ${process.env.DEVNET_USER_INDEX}`) + assert(userIndex >= 0 && userIndex < info.users.length, `DEVNET_USER_INDEX=${userIndex} out of range`) + + const devnetProvider = info.pdp_sps.find((provider) => provider.is_approved) + assert(devnetProvider != null, 'No approved PDP service provider found in devnet-info.json') + + const user = info.users[userIndex] + const chain = toChain(devnetInfo) + if (process.env.RPC_URL) { + chain.rpcUrls = { + ...chain.rpcUrls, + default: { http: [process.env.RPC_URL] }, + public: { http: [process.env.RPC_URL] }, + } + } + + const account = privateKeyToAccount(user.private_key_hex) + assert( + account.address.toLowerCase() === user.evm_addr.toLowerCase(), + `Derived address ${account.address} does not match ${user.name} address ${user.evm_addr}` + ) + + const source = 'foc-devnet-smoke' + const smokeId = process.env.CREATE_DATASET_SMOKE_ID || `cds-${Date.now().toString(36)}` + const metadata = { smoke: smokeId, source } + + console.log(`Devnet run: ${info.run_id}`) + console.log(`User: ${user.name} (${account.address})`) + console.log(`Provider: ${devnetProvider.provider_id} (${devnetProvider.pdp_service_url})`) + console.log(`Smoke metadata: smoke=${smokeId}, source=${source}`) + + const synapse = Synapse.create({ + chain, + transport: viemHttp(), + account, + source, + }) + + const provider = await synapse.providers.getProvider({ providerId: BigInt(devnetProvider.provider_id) }) + assert(provider != null, `Provider ${devnetProvider.provider_id} not found in registry`) + + const preExisting = await synapse.storage.findDataSets() + assert( + !preExisting.some((dataSet) => dataSet.metadata?.smoke === smokeId), + `Smoke metadata ${smokeId} already exists before createDataSet` + ) + + const context = await synapse.storage.createContext({ + providerId: provider.id, + metadata, + withCDN: false, + }) + assert(context.dataSetId == null, `Expected unique metadata to create a new data set, got ${context.dataSetId}`) + + await prepareWithPlainErc20(synapse, context) + + const create = await SP.createDataSet(synapse.client, { + cdn: false, + payee: provider.serviceProvider, + payer: account.address, + serviceURL: devnetProvider.pdp_service_url, + recordKeeper: chain.contracts.fwss.address, + metadata, + }) + console.log(`createDataSet tx submitted: ${create.txHash}`) + + const confirmed = await SP.waitForCreateDataSet({ + statusUrl: create.statusUrl, + timeout: 180000, + }) + assert(confirmed.dataSetId > 0n, `Expected positive dataSetId, got ${confirmed.dataSetId}`) + console.log(`createDataSet confirmed: dataSetId=${confirmed.dataSetId}`) + + const dataSet = await waitForDataSet(synapse, confirmed.dataSetId) + assert(dataSet.isLive === true, `Data set ${confirmed.dataSetId} is not live`) + assert(dataSet.isManaged === true, `Data set ${confirmed.dataSetId} is not managed by FWSS`) + assert(dataSet.providerId === provider.id, `Data set providerId ${dataSet.providerId} != ${provider.id}`) + assert( + dataSet.payer.toLowerCase() === account.address.toLowerCase(), + `Data set payer ${dataSet.payer} != ${account.address}` + ) + assert( + dataSet.payee.toLowerCase() === provider.serviceProvider.toLowerCase(), + `Data set payee ${dataSet.payee} != ${provider.serviceProvider}` + ) + assert(dataSet.metadata?.smoke === smokeId, `Data set smoke metadata missing or wrong: ${dataSet.metadata?.smoke}`) + assert(dataSet.metadata?.source === source, `Data set source metadata missing or wrong: ${dataSet.metadata?.source}`) + + console.log(`Verified createDataSet smoke dataSetId=${confirmed.dataSetId}`) +} + +main().catch((error) => { + console.error(error) + process.exit(1) +}) diff --git a/scenarios/run.py b/scenarios/run.py index 8a366a6..d0a23c5 100755 --- a/scenarios/run.py +++ b/scenarios/run.py @@ -26,6 +26,7 @@ ORDER = [ ("test_containers", 5), ("test_basic_balances", 10), + ("test_create_dataset_smoke", 300), ("test_storage_e2e", 200), ("test_multi_copy_upload", 600), ("test_caching_subsystem", 200), diff --git a/scenarios/synapse.py b/scenarios/synapse.py index abdca53..d6b96e9 100644 --- a/scenarios/synapse.py +++ b/scenarios/synapse.py @@ -4,9 +4,9 @@ from __future__ import annotations import os +import json import subprocess import time -import json from pathlib import Path from scenarios.dependencies import component @@ -104,19 +104,34 @@ def clone_and_build(tmp_dir: Path) -> Path | None: return sdk_dir -def upload_file(sdk_dir: Path, filepath: str, label: str): - """Upload a single file via example-storage-e2e.js.""" - env = {**os.environ, "NETWORK": "devnet"} - cmd = ["node", "utils/example-storage-e2e.js", str(filepath)] +def run_node_script( + sdk_dir: Path, + script_path: Path, + label: str, + args: list[str] | None = None, + env: dict | None = None, + timeout: int | None = None, +): + """Run a Node script in synapse-sdk with retry handling for transient state forks.""" + script_arg = str(script_path) + if script_path.is_absolute(): + try: + script_arg = str(script_path.relative_to(sdk_dir)) + except ValueError: + pass + + cmd = ["node", script_arg, *(args or [])] + process_env = {**os.environ, **(env or {})} max_attempts = len(UPLOAD_RETRY_DELAYS_SECS) + 1 for attempt in range(1, max_attempts + 1): result = subprocess.run( cmd, cwd=str(sdk_dir), - env=env, + env=process_env, text=True, capture_output=True, + timeout=timeout, ) details = "\n".join( part for part in (result.stderr.strip(), result.stdout.strip()) if part @@ -136,3 +151,15 @@ def upload_file(sdk_dir: Path, filepath: str, label: str): f"retrying in {delay}s (attempt {attempt}/{max_attempts})" ) time.sleep(delay) + + +def upload_file(sdk_dir: Path, filepath: str, label: str): + """Upload a single file via example-storage-e2e.js.""" + env = {**os.environ, "NETWORK": "devnet"} + run_node_script( + sdk_dir, + Path("utils/example-storage-e2e.js"), + label, + args=[str(filepath)], + env=env, + ) diff --git a/scenarios/test_create_dataset_smoke.py b/scenarios/test_create_dataset_smoke.py new file mode 100644 index 0000000..96dd1e8 --- /dev/null +++ b/scenarios/test_create_dataset_smoke.py @@ -0,0 +1,44 @@ +#!/usr/bin/env python3 +"""Smoke test for the PDP/FWSS createDataSet path via synapse-sdk.""" + +import os, sys # noqa: E401 + +sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) + +import shutil +import tempfile +from pathlib import Path + +from scenarios.helpers import assert_ok, info +from scenarios.synapse import clone_and_build, run_node_script + +SMOKE_TIMEOUT_SECS = 280 +SMOKE_USER_INDEX = "1" # USER_2; USER_1 is used by existing storage scenarios. +SMOKE_SCRIPT_SOURCE = Path(__file__).with_name("create-dataset-smoke.mjs") + + +def run(): + assert_ok("command -v git", "git is installed") + assert_ok("command -v node", "node is installed") + assert_ok("command -v pnpm", "pnpm is installed") + + with tempfile.TemporaryDirectory(prefix="synapse-sdk-createdataset-") as tmp: + sdk_dir = clone_and_build(Path(tmp)) + if not sdk_dir: + return + + script_path = sdk_dir / "utils" / SMOKE_SCRIPT_SOURCE.name + shutil.copyfile(SMOKE_SCRIPT_SOURCE, script_path) + + info("Running createDataSet smoke script against devnet") + run_node_script( + sdk_dir, + script_path, + "createDataSet smoke test", + env={"DEVNET_USER_INDEX": SMOKE_USER_INDEX}, + timeout=SMOKE_TIMEOUT_SECS, + ) + + +if __name__ == "__main__": + run() diff --git a/scripts/tests/test_scenario_dependencies.py b/scripts/tests/test_scenario_dependencies.py index 4f90e2a..78af245 100644 --- a/scripts/tests/test_scenario_dependencies.py +++ b/scripts/tests/test_scenario_dependencies.py @@ -1,10 +1,11 @@ import tempfile import unittest +import subprocess from pathlib import Path from unittest.mock import patch from scenarios.dependencies import format_markdown_table -from scenarios.synapse import clone_and_build +from scenarios.synapse import clone_and_build, run_node_script from scenarios.test_multi_copy_upload import setup_filecoin_pin @@ -77,6 +78,56 @@ def fake_run_cmd(command, **_kwargs): self.assertNotIn(["pnpm", "pkg", "set"], [command[:3] for command in commands]) self.assertIn(' "nanoid": "3.3.13"', workspace_text) + @patch("scenarios.synapse.ok") + @patch("scenarios.synapse.info") + @patch("scenarios.synapse.subprocess.run") + def test_run_node_script_uses_sdk_cwd_and_env(self, run, _info, ok): + run.return_value = subprocess.CompletedProcess( + ["node", "smoke.mjs"], 0, stdout="done\n", stderr="" + ) + with tempfile.TemporaryDirectory() as directory: + sdk_dir = Path(directory) + script = sdk_dir / "smoke.mjs" + run_node_script( + sdk_dir, + script, + "run smoke", + args=["random_file"], + env={"DEVNET_USER_INDEX": "1"}, + timeout=30, + ) + + kwargs = run.call_args.kwargs + self.assertEqual(run.call_args.args[0], ["node", "smoke.mjs", "random_file"]) + self.assertEqual(kwargs["cwd"], str(sdk_dir)) + self.assertEqual(kwargs["env"]["DEVNET_USER_INDEX"], "1") + self.assertEqual(kwargs["timeout"], 30) + ok.assert_called_once_with("run smoke") + + @patch("scenarios.synapse.time.sleep") + @patch("scenarios.synapse.ok") + @patch("scenarios.synapse.info") + @patch("scenarios.synapse.subprocess.run") + def test_run_node_script_retries_state_fork_error(self, run, _info, ok, sleep): + run.side_effect = [ + subprocess.CompletedProcess( + ["node", "smoke.mjs"], + 1, + stdout="", + stderr="refusing explicit call due to state fork at epoch 42", + ), + subprocess.CompletedProcess( + ["node", "smoke.mjs"], 0, stdout="done\n", stderr="" + ), + ] + + with tempfile.TemporaryDirectory() as directory: + run_node_script(Path(directory), Path("smoke.mjs"), "run smoke") + + self.assertEqual(run.call_count, 2) + sleep.assert_called_once_with(5) + ok.assert_called_once_with("run smoke") + @patch("scenarios.test_multi_copy_upload.run_cmd", return_value=True) @patch( "scenarios.test_multi_copy_upload.component",