diff --git a/Cargo.lock b/Cargo.lock index e0a7dcdf..dbc5e8ba 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4116,10 +4116,12 @@ dependencies = [ "starstream-compiler", "starstream-runtime-next", "starstream-to-wasm", + "starstream-types", "thiserror 2.0.17", "tokio", "tokio-util", "tracing", + "wasm-encoder 0.254.0", "wasm-tokio", "wasmparser 0.254.0", "wasmprinter 0.254.0", @@ -4150,6 +4152,7 @@ dependencies = [ "starstream-ledger", "starstream-runtime-next", "starstream-to-wasm", + "starstream-types", "tempfile", "tokio", "tokio-util", diff --git a/starstream-cli/src/run.rs b/starstream-cli/src/run.rs index 74fc5f26..968fb901 100644 --- a/starstream-cli/src/run.rs +++ b/starstream-cli/src/run.rs @@ -13,7 +13,7 @@ use tokio::fs; use tracing::{debug, info, instrument}; use wasmtime::component::{Component, Resource, ResourceTable, Val}; use wasmtime::error::Context as _; -use wasmtime::{AsContextMut as _, Store, StoreContextMut, ensure}; +use wasmtime::{AsContextMut as _, Store, StoreContextMut, bail, ensure}; /// Run a coordination script exported by a Wasm component #[derive(Args, Debug)] @@ -228,6 +228,25 @@ async fn exec( .call_coordination_script(&mut store, &script, ¶ms, &mut results) .await?; debug!(outputs = store.data().outputs.len(), "script returned"); + 'outer: for result in &mut results { + if let &mut Val::Resource(utxo) = result { + let utxo = utxo + .try_into_resource(&mut store) + .context("result resource is not a UTXO")?; + let utxo: &Utxo>> = store + .data() + .table + .get(&utxo) + .context("result UTXO not found")?; + for (i, out) in zip(0.., &store.data().outputs) { + if utxo.resource() == out.resource() { + *result = Val::U32(i); + continue 'outer; + } + } + bail!("failed to identify result resource"); + } + } let results = Val::Tuple(results) .to_wave() .context("failed to encode results")?; diff --git a/starstream-ledger-cli/Cargo.toml b/starstream-ledger-cli/Cargo.toml index 11ca7b5f..0e5e652e 100644 --- a/starstream-ledger-cli/Cargo.toml +++ b/starstream-ledger-cli/Cargo.toml @@ -36,6 +36,7 @@ hyper-util = { workspace = true, features = [ rand_core = { workspace = true, features = ["getrandom", "std"] } sha2 = { workspace = true } starstream-ledger = { workspace = true, features = ["client"] } +starstream-runtime-next = { workspace = true } toml = { workspace = true, features = ["display", "serde"] } tokio = { workspace = true, features = [ "fs", @@ -67,8 +68,8 @@ zeroize = { workspace = true } minicbor = { workspace = true, features = ["std"] } starstream-compiler = { workspace = true } starstream-ledger = { workspace = true, features = ["server"] } -starstream-runtime-next = { workspace = true } starstream-to-wasm = { workspace = true } +starstream-types = { workspace = true } tempfile = { workspace = true } tokio = { workspace = true, features = ["fs", "net", "process"] } toml = { workspace = true, features = ["parse"] } diff --git a/starstream-ledger-cli/src/main.rs b/starstream-ledger-cli/src/main.rs index 0afe8dad..f75ca670 100644 --- a/starstream-ledger-cli/src/main.rs +++ b/starstream-ledger-cli/src/main.rs @@ -5,6 +5,7 @@ use core::pin::pin; use std::collections::HashMap; use std::path::{Path, PathBuf}; +use std::sync::Arc; use anyhow::{Context as _, bail, ensure}; use bytes::{Bytes, BytesMut}; @@ -16,16 +17,19 @@ use hyper_util::rt::TokioExecutor; use rand_core::OsRng; use sha2::{Digest as _, Sha256}; use starstream_ledger::client::http::ClientBuilder; +use starstream_ledger::client::runtime::UtxoCtx; use starstream_ledger::client::runtime::{ - Client, call_coordination_script, compile_component, new_contract, + Client, Ctx, call_coordination_script, compile_component, new_contract, }; use starstream_ledger::client::{decode_transaction, encode_transaction}; use starstream_ledger::{TransactionInput, TransactionOutput, encode_digest}; +use starstream_runtime_next::Utxo; use tokio::fs; use tokio::io::{AsyncRead, AsyncWriteExt as _, stdout}; use tokio_util::codec::Encoder as _; use tracing::info; use wasm_wave::wasm::WasmFunc as _; +use wasmtime::Store; use wasmtime::component::{Val, types}; use zeroize::Zeroizing; @@ -69,6 +73,9 @@ enum Command { path: PathBuf, }, + /// Print the genesis outputs. + Genesis, + /// Manage signing keys. #[command(subcommand)] Key(KeyCommand), @@ -193,6 +200,12 @@ enum KeyCommand { #[derive(Debug, Subcommand)] enum TransactionCommand { + /// Get a transaction from the ledger. + Get { + /// Digest of the transaction, either as multibase multihash or `sha256:HEX`. + #[arg(value_parser = parse_digest)] + digest: [u8; 32], + }, /// Print a transaction written by `contract script call --output-transaction`. Show { /// Path to the encoded transaction. @@ -210,7 +223,7 @@ enum UtxoCommand { transaction: Option<[u8; 32]>, /// Index of the UTXO in the transaction outputs. - index: usize, + index: u32, /// Method to call. method: Box, @@ -359,6 +372,16 @@ async fn main() -> anyhow::Result<()> { .await .context("failed to write digest to stdout") } + Command::Genesis => { + let outputs = client.get_genesis().await?; + let outputs = toml::Value::try_from(outputs).context("failed to encode TOML")?; + let genesis = toml::Table::from_iter([("outputs".to_string(), outputs)]); + let genesis = toml::to_string_pretty(&genesis).context("failed to encode TOML")?; + stdout() + .write_all(genesis.as_bytes()) + .await + .context("failed to write genesis to stdout") + } Command::Contract(ContractCommand::Publish { signing: SigningArgs { key, nonce }, wasm, @@ -432,7 +455,10 @@ async fn main() -> anyhow::Result<()> { ensure!(args.next().is_none(), "trailing arguments"); let mut results = vec![Val::Bool(false); ty.results().len()]; + let mut utxos = Vec::default(); + let mut store = Store::new(client.engine(), Ctx::default()); let tx = call_coordination_script( + &mut store, &imports, client.wizer(), &contract, @@ -441,8 +467,29 @@ async fn main() -> anyhow::Result<()> { &mut contracts, params, &mut results, + &mut utxos, ) .await?; + 'outer: for result in &mut results { + if let &mut Val::Resource(utxo) = result { + let utxo = utxo + .try_into_resource(&mut store) + .map_err(anyhow::Error::from) + .context("result resource is not a UTXO")?; + let utxo: &Utxo>> = store + .data() + .table + .get(&utxo) + .context("result UTXO not found")?; + for (i, out) in zip(0.., &utxos) { + if utxo.resource() == out.resource() { + *result = Val::U32(i); + continue 'outer; + } + } + bail!("failed to identify result resource"); + } + } let mut results = wasm_wave::to_string(&Val::Tuple(results)) .context("failed to encode result tuple")?; if let Some(path) = output_transaction { @@ -471,6 +518,14 @@ async fn main() -> anyhow::Result<()> { .await .context("failed to write signing key to stdout") } + Command::Transaction(TransactionCommand::Get { digest }) => { + let tx = client.get_transaction(digest).await?; + let tx = toml::to_string_pretty(&tx).context("failed to encode TOML")?; + stdout() + .write_all(tx.as_bytes()) + .await + .context("failed to write transaction to stdout") + } Command::Transaction(TransactionCommand::Show { path }) => { let buf = fs::read(&path) .await @@ -488,17 +543,18 @@ async fn main() -> anyhow::Result<()> { method, args, }) => { - let TransactionOutput { - instance, - methods, - storage, - wasm, - .. - } = if let Some(transaction) = transaction { - client.get_transaction_utxo(transaction, index).await? + let transaction = if let Some(transaction) = transaction { + encode_digest(&transaction).into() } else { - client.get_genesis_utxo(index).await? + Box::default() }; + let input = TransactionInput { transaction, index }; + let TransactionOutput { + contract, instance, .. + } = client.get_input_utxo(&input).await?; + let digest = starstream_ledger::parse_digest(&contract) + .with_context(|| format!("failed to parse `{contract}` as multibase multihash"))?; + let wasm = client.get_contract_wasm(digest).await?; let (resolve, world) = decode_component(&wasm)?; let world = &resolve.worlds[world]; let ty = world @@ -529,10 +585,7 @@ async fn main() -> anyhow::Result<()> { let ty = wasm_wave::value::resolve_wit_func_type(&resolve, &ty) .context("failed to resolve method type")?; let args = encode_args(&ty, args)?; - let digest = Sha256::digest(&wasm).into(); - let rx = client - .call_utxo_method(&digest, &instance, &method, &methods, &storage, &args) - .await?; + let rx = client.call_utxo_method(&input, &method, &args).await?; write_results(&ty, rx).await } } diff --git a/starstream-ledger-cli/tests/cli.rs b/starstream-ledger-cli/tests/cli.rs index 0ac68626..111ccf71 100644 --- a/starstream-ledger-cli/tests/cli.rs +++ b/starstream-ledger-cli/tests/cli.rs @@ -83,16 +83,16 @@ async fn cli() { .expect("failed to handle HTTP"); let ledger = tokio::spawn(ledger); - let wasm = NamedTempFile::new().unwrap(); - fs::write(&wasm, &*SCORE_WASM) + let score_wasm = NamedTempFile::new().unwrap(); + fs::write(&score_wasm, &*SCORE_WASM) .await - .with_context(|| format!("failed to write Wasm to `{}`", wasm.path().display())) + .with_context(|| format!("failed to write Wasm to `{}`", score_wasm.path().display())) .unwrap(); - let digest = run_cli(["digest", &wasm.path().to_string_lossy()]) + let score_digest = run_cli(["digest", &score_wasm.path().to_string_lossy()]) .await .unwrap(); - let digest = str::from_utf8(&digest) + let score_digest = str::from_utf8(&score_digest) .expect("contract digest is not valid UTF-8") .trim_end(); @@ -109,15 +109,15 @@ async fn cli() { "--output-transaction", &tx_file.path().to_string_lossy(), "--import", - &wasm.path().to_string_lossy(), - digest, + &score_wasm.path().to_string_lossy(), + score_digest, "example", ]) .await .unwrap(); assert_eq!(stdout, b"()\n"); let tx = fs::read(&tx_file).await.unwrap(); - let envelope = assert_score_transaction(&tx, digest); + let envelope = assert_score_transaction(&tx, score_digest); let stdout = run_cli(["transaction", "show", &tx_file.path().to_string_lossy()]) .await @@ -170,7 +170,7 @@ async fn cli() { NETWORK, "--nonce", "1", - &wasm.path().to_string_lossy(), + &score_wasm.path().to_string_lossy(), ]) .await .unwrap(); @@ -200,7 +200,39 @@ async fn cli() { .unwrap(); assert_eq!(stdout, b"()\n"); let tx = fs::read(&tx_file).await.unwrap(); - assert_score_transaction(&tx, digest); + assert_score_transaction(&tx, score_digest); + + let require_hash_preimage_wasm = NamedTempFile::new().unwrap(); + fs::write(&require_hash_preimage_wasm, &*REQUIRE_HASH_PREIMAGE_WASM) + .await + .with_context(|| { + format!( + "failed to write Wasm to `{}`", + require_hash_preimage_wasm.path().display() + ) + }) + .unwrap(); + let stdout = run_cli([ + "--url", + &format!("http://{addr}"), + "contract", + "script", + "call", + "--network", + NETWORK, + "--simulate", + &format!( + "sha256:{}", + hex::encode(Sha256::digest(&*REQUIRE_HASH_PREIMAGE_WASM)) + ), + "--import", + &require_hash_preimage_wasm.path().to_string_lossy(), + "create-hash", + "42", + ]) + .await + .unwrap(); + assert_eq!(stdout, b"(0)\n"); shutdown.notify_one(); ledger.await.expect("ledger task panicked"); diff --git a/starstream-ledger-node/src/main.rs b/starstream-ledger-node/src/main.rs index 2bbcc311..55fd5fd7 100644 --- a/starstream-ledger-node/src/main.rs +++ b/starstream-ledger-node/src/main.rs @@ -15,7 +15,7 @@ use ed25519_dalek::VerifyingKey; use serde::Deserialize; use sha2::{Digest as _, Sha256}; use starstream_ledger::client::runtime::{ - Client, Contract, call_coordination_script, compile_component, new_contract, + Client, Contract, Ctx, call_coordination_script, compile_component, new_contract, }; use starstream_ledger::server::Ledger; use starstream_ledger::{TransactionInput, TransactionOutput, encode_digest, parse_digest}; @@ -23,6 +23,7 @@ use tokio::fs; use tokio::signal; use tokio::time::timeout; use tracing::{error, info, warn}; +use wasmtime::Store; use wasmtime::component::Val; use wasmtime_wizer::Wizer; @@ -144,12 +145,13 @@ async fn build_genesis( ) .await?; let contract = Contract { - contract, + contract: Some(contract), wasm: wasm.clone(), }; contracts.insert(digest, contract.clone()); contract }; + let contract = contract.context("contract was not compiled")?; let script = contract.get_coordination_script(&script)?; let ty = script.ty(); let mut params = Vec::with_capacity(ty.params().len()); @@ -165,6 +167,7 @@ async fn build_genesis( ensure!(args.next().is_none(), "trailing arguments"); let mut results = vec![Val::Bool(false); ty.results().len()]; let tx = call_coordination_script( + &mut Store::new(engine, Ctx::default()), &GenesisClient(&imported), &wizer, &contract, @@ -173,6 +176,7 @@ async fn build_genesis( &mut contracts, params, &mut results, + &mut Vec::default(), ) .await?; outputs.extend(tx.outputs); diff --git a/starstream-ledger/Cargo.toml b/starstream-ledger/Cargo.toml index 6faa274e..4047c27e 100644 --- a/starstream-ledger/Cargo.toml +++ b/starstream-ledger/Cargo.toml @@ -27,7 +27,6 @@ server = [ "dep:hyper", "dep:tokio", "dep:wasm-tokio", - "dep:wasmparser", "dep:wasmtime", "hyper-util/server", "hyper-util/server-auto", @@ -66,7 +65,8 @@ tokio = { workspace = true, optional = true, features = [ tokio-util = { workspace = true, features = ["codec"] } tracing = { workspace = true, features = ["attributes"] } wasm-tokio = { workspace = true, optional = true } -wasmparser = { workspace = true, optional = true } +wasm-encoder = { workspace = true, features = ["wasmparser"] } +wasmparser = { workspace = true } wasmprinter = { workspace = true } wasmtime = { workspace = true, optional = true, features = [ "anyhow", @@ -85,6 +85,7 @@ wrpc-transport = { workspace = true } [dev-dependencies] starstream-compiler = { workspace = true } starstream-to-wasm = { workspace = true } +starstream-types = { workspace = true } tokio = { workspace = true, features = [ "io-util", "macros", diff --git a/starstream-ledger/src/client/http.rs b/starstream-ledger/src/client/http.rs index 6b1f7945..094f8917 100644 --- a/starstream-ledger/src/client/http.rs +++ b/starstream-ledger/src/client/http.rs @@ -1,4 +1,5 @@ -use std::collections::{BTreeSet, HashMap}; +use std::collections::HashMap; +use std::sync::Arc; use anyhow::{Context as _, ensure}; use bytes::{Bytes, BytesMut}; @@ -12,19 +13,22 @@ use http_body_util::{BodyExt as _, Full}; use hyper_util::client::legacy::connect::Connect; use mediatype::MediaType; use sha2::{Digest as _, Sha256}; -use starstream_runtime_next::CoordinationScriptExport; +use starstream_runtime_next::{CoordinationScriptExport, Utxo}; use tokio_util::codec::Encoder as _; use tracing::{instrument, warn}; +use wasm_tokio::cm::OptionEncoder; use wasm_tokio::{CoreNameEncoder, Leb128Encoder}; +use wasmtime::Store; use wasmtime::component::Val; use wasmtime_wizer::Wizer; use wrpc_transport::Invoke as _; -use crate::client::runtime::{Contract, Ctx, call_coordination_script}; +use crate::client::runtime::{Contract, Ctx, UtxoCtx, call_coordination_script}; use crate::client::{ CoordinationScriptArg, bindings, build_fund_envelope, build_publish_envelope, - build_sign_envelope, encode_transaction, utxo_instance, + build_sign_envelope, encode_transaction, }; +use crate::wrpc::LEDGER_PACKAGE; use crate::{ APPLICATION_COSE, APPLICATION_WASM, Envelope, EnvelopeContext, Fund, Publish, Transaction, TransactionInput, TransactionOutput, encode_digest, parse_digest, @@ -314,16 +318,20 @@ where /// Contracts in `imports` are used to resolve imports instead of the ledger. /// `wasm` must be equal to original component bytes. #[instrument(skip_all)] + #[allow(clippy::too_many_arguments)] pub async fn call_coordination_script( &self, + store: &mut Store, contract: &starstream_runtime_next::Contract, wasm: &[u8], export: &CoordinationScriptExport, imports: &mut HashMap<[u8; 32], Contract>, args: impl IntoIterator, results: &mut [Val], + utxos: &mut Vec>>>, ) -> anyhow::Result { let tx = call_coordination_script( + store, self, &self.wizer, contract, @@ -332,50 +340,44 @@ where imports, args, results, + utxos, ) .await?; Ok(tx) } /// Call the method `name` exported by the UTXO - /// identified by `digest` with encoded `args`. - /// - /// `methods` is the set of method hashes the UTXO implements and - /// `storage` its encoded storage record, both as found in the - /// [`TransactionOutput`] the UTXO was created by. + /// created by the transaction output `input` with encoded `args`. #[instrument(skip_all)] pub async fn call_utxo_method( &self, - digest: &[u8; 32], - instance: &str, + TransactionInput { transaction, index }: &TransactionInput, name: &str, - methods: &BTreeSet<(u64, u64, u64, u64)>, - storage: &[u8], args: &[u8], ) -> anyhow::Result { let cx = wrpc_context(&self.api_base)?; - let mut params = BytesMut::with_capacity( - 5 + instance.len() + 5 + methods.len() * 40 + storage.len() + args.len(), - ); - CoreNameEncoder - .encode(instance, &mut params) - .context("failed to encode instance name")?; - let n = u32::try_from(methods.len()).context("method set length does not fit in u32")?; + let mut params = BytesMut::with_capacity(1 + 5 + transaction.len() + 5 + args.len()); + let transaction = if transaction.is_empty() { + None + } else { + Some(transaction) + }; + OptionEncoder(CoreNameEncoder) + .encode(transaction, &mut params) + .context("failed to encode transaction digest")?; Leb128Encoder - .encode(n, &mut params) - .context("failed to encode method set length")?; - for &(a, b, c, d) in methods { - for v in [a, b, c, d] { - Leb128Encoder - .encode(v, &mut params) - .context("failed to encode method hash")?; - } - } - params.extend_from_slice(storage); + .encode(*index, &mut params) + .context("failed to encode output index")?; params.extend_from_slice(args); let (tx, rx) = self .wrpc - .invoke(cx, &utxo_instance(digest), name, params.freeze(), [[]]) + .invoke( + cx, + &format!("{LEDGER_PACKAGE}/utxo"), + name, + params.freeze(), + [[]], + ) .await?; drop(tx); Ok(rx) diff --git a/starstream-ledger/src/client/mod.rs b/starstream-ledger/src/client/mod.rs index 61cdf9e3..809e0bbf 100644 --- a/starstream-ledger/src/client/mod.rs +++ b/starstream-ledger/src/client/mod.rs @@ -7,10 +7,7 @@ use coset::{ use ed25519_dalek::{Signer as _, SigningKey}; use wasmtime::component::Val; -use crate::wrpc::UTXO_PACKAGE; -use crate::{ - Envelope, EnvelopeContext, Fund, Publish, Transaction, TransactionInput, encode_digest, -}; +use crate::{Envelope, EnvelopeContext, Fund, Publish, Transaction, TransactionInput}; pub mod http; pub mod runtime; @@ -20,11 +17,6 @@ pub mod bindings { wit_bindgen_wrpc::generate!(); } -/// The wRPC instance name of the UTXO identified by `digest`. -fn utxo_instance(digest: &[u8; 32]) -> String { - format!("{UTXO_PACKAGE}/{}", encode_digest(digest)) -} - /// Build a signed `COSE_Sign` envelope. fn build_sign_envelope(key: SigningKey, payload: impl Into>) -> anyhow::Result> { let protected = coset::HeaderBuilder::new() diff --git a/starstream-ledger/src/client/runtime.rs b/starstream-ledger/src/client/runtime.rs index eecef432..800c654f 100644 --- a/starstream-ledger/src/client/runtime.rs +++ b/starstream-ledger/src/client/runtime.rs @@ -15,10 +15,11 @@ use tokio_util::codec::Encoder as _; use tracing::error; use wasmtime::component::{Component, Resource, ResourceTable, Type, Val}; use wasmtime::error::Context as _; -use wasmtime::{AsContextMut as _, Engine, StoreContextMut, bail, ensure, format_err}; +use wasmtime::{AsContextMut as _, Engine, Store, StoreContextMut, bail, ensure, format_err}; use wasmtime_wizer::{WasmtimeWizerComponent, Wizer}; use crate::client::CoordinationScriptArg; +use crate::runtime::{apply_state, parse_state}; use crate::wrpc::codec::{ValEncoder, read_value}; use crate::{ Transaction, TransactionEvent, TransactionInput, TransactionOutput, encode_digest, parse_digest, @@ -39,7 +40,7 @@ pub fn compile_component( #[derive(Clone)] pub struct Contract { - pub contract: starstream_runtime_next::Contract, + pub contract: Option>, pub wasm: Bytes, } @@ -70,7 +71,10 @@ pub async fn new_contract( error!(external_id, "unresolved contract import"); format!("contract identified by `external-id` `{external_id}` not found") })?; - Ok(contract.contract.clone()) + contract.contract.clone().with_context(|| { + error!(external_id, "uncompiled contract import"); + format!("contract identified by `external-id` `{external_id}` was not compiled") + }) } } @@ -89,14 +93,20 @@ pub async fn new_contract( let digest = parse_digest(external_id).with_context(|| { format!("failed to parse `external-id` `{external_id}` as multibase multihash") })?; - if imports.contains_key(&digest) { - continue; - } - let wasm = client - .get_contract_wasm(digest) - .await - .map_err(wasmtime::Error::from_anyhow)?; - let component = compile_component(engine, wizer, &wasm)?; + let wasm = match imports.get(&digest) { + Some(Contract { + contract: Some(..), .. + }) => continue, + Some(Contract { + contract: None, + wasm, + }) => wasm.clone(), + None => client + .get_contract_wasm(digest) + .await + .map_err(wasmtime::Error::from_anyhow)?, + }; + let component = compile_component(engine, wizer, wasm.as_ref())?; let contract = Box::pin(new_contract( client, wizer, @@ -105,13 +115,20 @@ pub async fn new_contract( imports, )) .await?; - imports.insert(digest, Contract { contract, wasm }); + imports.insert( + digest, + Contract { + contract: Some(contract), + wasm, + }, + ); } starstream_runtime_next::Contract::new(component, external_id, ContractLookup(imports)) } #[allow(clippy::too_many_arguments)] pub async fn call_coordination_script( + store: &mut Store, client: &(impl Client + ?Sized), wizer: &Wizer, contract: &starstream_runtime_next::Contract, @@ -120,10 +137,10 @@ pub async fn call_coordination_script( imports: &mut HashMap<[u8; 32], Contract>, args: impl IntoIterator, results: &mut [Val], + utxos: &mut Vec>>>, ) -> wasmtime::Result { let engine = contract.component().engine(); - let mut store = wasmtime::Store::new(engine, Ctx::default()); - let instance = contract.instantiate(&mut store).await?; + let instance = contract.instantiate(&mut *store).await?; let digest = Sha256::digest(wasm).into(); let mut inputs = Vec::default(); @@ -143,33 +160,40 @@ pub async fn call_coordination_script( let utxo_contract_digest = parse_digest(&utxo.contract).with_context(|| { format!("failed to parse `{}` as multibase multihash", utxo.contract) })?; - let (external_id, contract) = if utxo_contract_digest == digest { - (None, contract.clone()) - } else if let Some(Contract { contract, .. }) = imports.get(&utxo_contract_digest) { - (Some(Arc::from(utxo.contract)), contract.clone()) + let (external_id, wasm) = if utxo_contract_digest == digest { + let wasm = + apply_state(wasm, &utxo.state).map_err(wasmtime::Error::from_anyhow)?; + (None, Bytes::from(wasm)) + } else if let Some(Contract { wasm, .. }) = imports.get(&utxo_contract_digest) { + let wasm = + apply_state(wasm, &utxo.state).map_err(wasmtime::Error::from_anyhow)?; + (Some(Arc::from(utxo.contract)), Bytes::from(wasm)) } else { let wasm = client .get_contract_wasm(utxo_contract_digest) .await .map_err(wasmtime::Error::from_anyhow)?; - let component = compile_component(engine, wizer, &wasm)?; - let contract = new_contract( - client, - wizer, - &component, - Some(&utxo.contract), - &mut *imports, - ) - .await?; + let wasm = + apply_state(&wasm, &utxo.state).map_err(wasmtime::Error::from_anyhow)?; + let wasm = Bytes::from(wasm); imports.insert( utxo_contract_digest, Contract { - contract: contract.clone(), - wasm, + contract: None, + wasm: wasm.clone(), }, ); - (Some(Arc::from(utxo.contract)), contract) + (Some(Arc::from(utxo.contract)), wasm) }; + let component = compile_component(engine, wizer, &wasm)?; + let contract = new_contract( + client, + wizer, + &component, + external_id.as_deref(), + &mut *imports, + ) + .await?; let utxo_export = contract.get_utxo(&utxo.instance)?; let storage_export = utxo_export.storage().context("UTXO has no storage")?; let mut storage = Val::Record(Vec::default()); @@ -189,11 +213,11 @@ pub async fn call_coordination_script( dropped: false, })); let cx_res = store.data_mut().table.push(Arc::clone(&cx))?; - let cx_res = cx_res.try_into_resource_any(&mut store)?; - let contract = contract.instantiate(&mut store).await?; + let cx_res = cx_res.try_into_resource_any(&mut *store)?; + let contract = contract.instantiate(&mut *store).await?; let utxo = contract .load_utxo( - &mut store, + &mut *store, &utxo_export, storage_export, cx, @@ -203,7 +227,7 @@ pub async fn call_coordination_script( let Ctx { table, outputs, .. } = store.data_mut(); outputs.push(utxo.clone()); let utxo = table.push(utxo)?; - let utxo = utxo.try_into_resource_any(&mut store)?; + let utxo = utxo.try_into_resource_any(&mut *store)?; inputs.push(input); Val::Resource(utxo) } @@ -212,7 +236,7 @@ pub async fn call_coordination_script( } ensure!(args.next().is_none(), "trailing arguments"); instance - .call_coordination_script(&mut store, export, ¶ms, results) + .call_coordination_script(&mut *store, export, ¶ms, results) .await?; let Ctx { outputs, events, .. @@ -228,7 +252,7 @@ pub async fn call_coordination_script( cx.clone() }; let mut instance = WasmtimeWizerComponent { - store: &mut store, + store: &mut *store, instance: utxo.instance(), }; let (contract, wasm) = if let Some(external_id) = cx.external_id.as_deref() { @@ -244,8 +268,12 @@ pub async fn call_coordination_script( let wasm = wizer.snapshot_component(&wizer_cx, &mut instance).await?; (encode_digest(&digest).into(), wasm) }; + let state = parse_state(&wasm) + .collect::>() + .map_err(wasmtime::Error::from_anyhow) + .context("failed to parse UTXO state")?; let storage = if let Some(export) = cx.export.storage() { - let storage = utxo.storage(export).call_get(&mut store).await?; + let storage = utxo.storage(export).call_get(&mut *store).await?; let mut buf = BytesMut::new(); ValEncoder::new(&Type::Record(export.ty().clone())) .encode(&Val::Record(storage), &mut buf) @@ -258,13 +286,13 @@ pub async fn call_coordination_script( for &(a, b, c, d) in &cx.methods { methods.insert((a, b, c, d)); } - // TODO: remove contract code from UTXO snapshot + utxos.push(utxo); tx_outputs.push(TransactionOutput { contract, instance: cx.instance.as_ref().into(), methods, storage, - wasm: wasm.into(), + state, }); } Ok(Transaction { diff --git a/starstream-ledger/src/lib.rs b/starstream-ledger/src/lib.rs index 622d2b33..b4e053a8 100644 --- a/starstream-ledger/src/lib.rs +++ b/starstream-ledger/src/lib.rs @@ -10,11 +10,14 @@ use minicbor::{Decode, Encode}; use serde::Serialize; use thiserror::Error; +use crate::runtime::ModuleState; + #[cfg(feature = "client")] pub mod client; #[cfg(feature = "server")] pub mod server; +pub mod runtime; pub mod wrpc; pub const FUND_CONTEXT: &str = "starstream:fund"; @@ -89,9 +92,8 @@ pub struct TransactionOutput { #[cbor(n(3), with = "minicbor::bytes")] #[serde(serialize_with = "serialize_bytes")] pub storage: Box<[u8]>, - #[cbor(n(4), with = "minicbor::bytes")] - #[serde(serialize_with = "serialize_wasm")] - pub wasm: Box<[u8]>, + #[n(4)] + pub state: Vec, } #[derive(Debug, Clone, Eq, PartialEq, Ord, PartialOrd, Encode, Decode, Serialize)] diff --git a/starstream-ledger/src/runtime.rs b/starstream-ledger/src/runtime.rs new file mode 100644 index 00000000..2aa2d9e4 --- /dev/null +++ b/starstream-ledger/src/runtime.rs @@ -0,0 +1,333 @@ +use core::slice; + +use anyhow::{Context as _, bail, ensure}; +use minicbor::{Decode, Encode}; +use serde::Serialize; +use wasm_encoder::reencode::{Reencode, ReencodeComponent}; +use wasmparser::{DataSectionReader, MemorySectionReader, Parser, Payload}; + +#[derive(Debug, Clone, Eq, PartialEq, Ord, PartialOrd, Encode, Decode, Serialize)] +pub enum GlobalValue { + #[n(0)] + I32(#[n(0)] i32), + #[n(1)] + I64(#[n(0)] i64), + #[n(2)] + F32(#[n(0)] u32), + #[n(3)] + F64(#[n(0)] u64), + #[n(4)] + V128(#[cbor(n(0), with = "minicbor::bytes")] [u8; 16]), +} + +#[derive(Debug, Clone, Eq, PartialEq, Ord, PartialOrd, Encode, Decode, Serialize)] +pub struct DataState { + #[n(0)] + pub memory_index: u32, + #[n(1)] + pub offset: u32, + #[cbor(n(2), with = "minicbor::bytes")] + #[serde(serialize_with = "crate::serialize_bytes")] + pub data: Vec, +} + +impl DataState { + fn encode(&self, section: &mut wasm_encoder::DataSection) { + section.active( + self.memory_index, + &wasm_encoder::ConstExpr::i32_const(self.offset.cast_signed()), + self.data.iter().copied(), + ); + } +} + +#[derive(Default, Debug, Clone, Eq, PartialEq, Ord, PartialOrd, Encode, Decode, Serialize)] +pub struct ModuleState { + #[n(0)] + pub memories: Vec, + #[n(1)] + pub globals: Vec, + #[n(2)] + pub data: Vec, +} + +/// Parse the state of every core module in `wasm`, in definition order. +/// This function performs no input validation and assumes Wasm to be valid. +pub fn parse_state(wasm: &[u8]) -> impl Iterator> { + Parser::new(0) + .parse_all(wasm) + .filter_map(|payload| match payload { + Ok(Payload::ModuleSection { + parser, + unchecked_range, + }) => Some(parse_module_section(parser, &wasm[unchecked_range])), + Err(err) => Some(Err(err.into())), + _ => None, + }) +} + +fn parse_module_section(parser: wasmparser::Parser, wasm: &[u8]) -> anyhow::Result { + let mut state = ModuleState::default(); + for payload in parser.parse_all(wasm) { + let payload = payload?; + match payload { + Payload::ImportSection(reader) => { + for import in reader.into_imports() { + let wasmparser::Import { module, name, ty } = import?; + if let wasmparser::TypeRef::Memory(..) | wasmparser::TypeRef::Global(..) = ty { + bail!( + "module `{module}` imports `{name}`, which is of unsupported type `{ty:?}`", + ) + } + } + } + Payload::MemorySection(reader) => { + for memory in reader.into_iter() { + let wasmparser::MemoryType { initial, .. } = memory?; + state.memories.push(initial); + } + } + Payload::GlobalSection(reader) => { + for global in reader.into_iter() { + let wasmparser::Global { + ty: wasmparser::GlobalType { content_type, .. }, + init_expr, + } = global?; + ensure!(content_type.is_defaultable()); + + let mut init_expr = init_expr.get_operators_reader(); + let op = init_expr.read()?; + ensure!( + init_expr.is_end_then_eof(), + "global initialization expression is not a single instruction" + ); + match op { + wasmparser::Operator::I32Const { value } => { + state.globals.push(GlobalValue::I32(value)) + } + wasmparser::Operator::I64Const { value } => { + state.globals.push(GlobalValue::I64(value)) + } + wasmparser::Operator::F32Const { value } => { + state.globals.push(GlobalValue::F32(value.bits())) + } + wasmparser::Operator::F64Const { value } => { + state.globals.push(GlobalValue::F64(value.bits())) + } + wasmparser::Operator::V128Const { value } => { + state.globals.push(GlobalValue::V128(*value.bytes())) + } + op => { + bail!("unexpected global initialization expression operator `{op:?}`") + } + } + } + } + Payload::DataSection(reader) => { + for data in reader.into_iter() { + let wasmparser::Data { + kind: + wasmparser::DataKind::Active { + memory_index, + offset_expr, + }, + data, + .. + } = data? + else { + continue; + }; + let mut offset_expr = offset_expr.get_operators_reader(); + let op = offset_expr.read()?; + ensure!( + offset_expr.is_end_then_eof(), + "active data section offset expression is not a single instruction" + ); + let offset = match op { + wasmparser::Operator::I32Const { value } => value, + op => { + bail!( + "unexpected active data section offset expression operator `{op:?}`" + ) + } + }; + state.data.push(DataState { + memory_index, + offset: offset.cast_unsigned(), + data: data.into(), + }) + } + } + _ => {} + } + } + Ok(state) +} + +type ReencodeError = wasm_encoder::reencode::Error; + +/// Apply component the state previously parsed via [parse_state]. +pub fn apply_state<'a>( + wasm: &[u8], + state: impl IntoIterator, +) -> anyhow::Result> { + let state = state.into_iter(); + let mut enc = ComponentEncoder { state }; + let mut component = wasm_encoder::Component::new(); + enc.parse_component(&mut component, Parser::new(0), wasm) + .map_err(|err| match err { + ReencodeError::UserError(err) => err, + ReencodeError::ParseError(err) => err.into(), + err => anyhow::Error::msg(err), + })?; + ensure!( + enc.state.next().is_none(), + "state defines more core modules than the contract" + ); + Ok(component.finish()) +} + +struct ComponentEncoder { + state: I, +} + +impl Reencode for ComponentEncoder { + type Error = anyhow::Error; +} + +impl<'a, I: Iterator> ReencodeComponent for ComponentEncoder { + fn parse_component_submodule( + &mut self, + component: &mut wasm_encoder::Component, + parser: Parser, + wasm: &[u8], + ) -> Result<(), ReencodeError> { + let state = self + .state + .next() + .context("state defines fewer core modules than the contract") + .map_err(ReencodeError::UserError)?; + let mut encoder = ModuleEncoder { + state, + memories: state.memories.iter(), + globals: state.globals.iter(), + data: Some(&state.data), + }; + let mut module = wasm_encoder::Module::new(); + wasm_encoder::reencode::utils::parse_core_module(&mut encoder, &mut module, parser, wasm)?; + encoder + .finish(&mut module) + .map_err(ReencodeError::UserError)?; + component.section(&wasm_encoder::ModuleSection(&module)); + Ok(()) + } +} + +struct ModuleEncoder<'a> { + state: &'a ModuleState, + memories: slice::Iter<'a, u64>, + globals: slice::Iter<'a, GlobalValue>, + data: Option<&'a [DataState]>, +} + +impl ModuleEncoder<'_> { + fn finish(mut self, module: &mut wasm_encoder::Module) -> anyhow::Result<()> { + ensure!( + self.memories.next().is_none(), + "state defines more memories than the contract module" + ); + ensure!( + self.globals.next().is_none(), + "state defines more globals than the contract module" + ); + if let Some(segments) = self.data + && !segments.is_empty() + { + let mut section = wasm_encoder::DataSection::new(); + for state in segments { + state.encode(&mut section); + } + module.section(§ion); + } + Ok(()) + } +} + +impl Reencode for ModuleEncoder<'_> { + type Error = anyhow::Error; + + fn parse_memory_section( + &mut self, + memories: &mut wasm_encoder::MemorySection, + section: MemorySectionReader<'_>, + ) -> Result<(), ReencodeError> { + for memory in section { + let ty = self.memory_type(memory?)?; + let minimum = self + .memories + .next() + .context("state defines fewer memories than the contract module") + .map_err(ReencodeError::UserError)?; + memories.memory(wasm_encoder::MemoryType { + minimum: *minimum, + ..ty + }); + } + Ok(()) + } + + fn parse_global( + &mut self, + globals: &mut wasm_encoder::GlobalSection, + wasmparser::Global { ty, .. }: wasmparser::Global<'_>, + ) -> Result<(), ReencodeError> { + let ty = self.global_type(ty)?; + let value = self + .globals + .next() + .context("state defines fewer globals than the contract module") + .map_err(ReencodeError::UserError)?; + let expr = match value { + GlobalValue::I32(v) => wasm_encoder::ConstExpr::i32_const(*v), + GlobalValue::I64(v) => wasm_encoder::ConstExpr::i64_const(*v), + GlobalValue::F32(v) => { + wasm_encoder::ConstExpr::f32_const(wasm_encoder::Ieee32::new(*v)) + } + GlobalValue::F64(v) => { + wasm_encoder::ConstExpr::f64_const(wasm_encoder::Ieee64::new(*v)) + } + GlobalValue::V128(v) => wasm_encoder::ConstExpr::v128_const(i128::from_le_bytes(*v)), + }; + globals.global(ty, &expr); + Ok(()) + } + + fn data_count(&mut self, count: u32) -> Result { + let n = u32::try_from(self.state.data.len()) + .context("state data segment count does not fit in u32") + .map_err(ReencodeError::UserError)?; + n.checked_add(count) + .context("merged data segment count does not fit in u32") + .map_err(ReencodeError::UserError) + } + + fn parse_data_section( + &mut self, + data: &mut wasm_encoder::DataSection, + section: DataSectionReader<'_>, + ) -> Result<(), ReencodeError> { + for seg in section { + let seg = seg?; + match seg.kind { + wasmparser::DataKind::Active { .. } => data.passive([]), + wasmparser::DataKind::Passive => data.passive(seg.data.iter().copied()), + }; + } + if let Some(segments) = self.data.take() { + for state in segments { + state.encode(data); + } + } + Ok(()) + } +} diff --git a/starstream-ledger/src/server/http/error.rs b/starstream-ledger/src/server/http/error.rs index 2c813a5c..4ec9f011 100644 --- a/starstream-ledger/src/server/http/error.rs +++ b/starstream-ledger/src/server/http/error.rs @@ -203,8 +203,6 @@ impl TransactionGetError { #[derive(Debug, Error)] pub enum GenesisGetError { - #[error("failed to encode genesis: {0}")] - Encoding(minicbor::encode::Error), #[error(transparent)] Http(http::Error), } @@ -212,7 +210,7 @@ pub enum GenesisGetError { impl GenesisGetError { pub fn http_status_code(&self) -> http::StatusCode { match self { - Self::Encoding(..) | Self::Http(..) => http::StatusCode::INTERNAL_SERVER_ERROR, + Self::Http(..) => http::StatusCode::INTERNAL_SERVER_ERROR, } } } @@ -278,18 +276,28 @@ pub enum RpcPostError { InstanceNotFound(String), #[error("function `{name}` not found in instance `{instance}`")] FunctionNotFound { instance: String, name: String }, - #[error("failed to parse utxo digest: {0}")] - UtxoDigestParsing(DigestParseError), + #[error("failed to parse transaction digest: {0}")] + TransactionDigestParsing(DigestParseError), + #[error("UTXO index does not fit in usize")] + UtxoIndexOverflow, #[error("UTXO not found")] UtxoNotFound, + #[error("failed to parse contract digest: {0}")] + ContractDigestParsing(DigestParseError), + #[error("contract not found")] + ContractNotFound, + #[error("failed to merge UTXO state into contract: {0:#}")] + StateMerge(anyhow::Error), + #[error("failed to decode UTXO storage: {0}")] + StorageDecoding(std::io::Error), #[error("UTXO instance `{instance}` not found: {source:#}")] UtxoInstanceNotFound { - instance: String, + instance: Box, source: wasmtime::Error, }, #[error("method `{name}` not found in UTXO instance `{instance}`: {source:#}")] UtxoMethodNotFound { - instance: String, + instance: Box, name: String, source: wasmtime::Error, }, @@ -315,16 +323,21 @@ impl RpcPostError { pub fn http_status_code(&self) -> http::StatusCode { match self { Self::Header(..) - | Self::UtxoDigestParsing(..) + | Self::TransactionDigestParsing(..) + | Self::UtxoIndexOverflow | Self::ParameterDecoding(..) | Self::UtxoStorageMissing | Self::ResourceTable(..) => http::StatusCode::BAD_REQUEST, Self::InstanceNotFound(..) | Self::FunctionNotFound { .. } | Self::UtxoNotFound + | Self::ContractNotFound | Self::UtxoInstanceNotFound { .. } | Self::UtxoMethodNotFound { .. } => http::StatusCode::NOT_FOUND, - Self::Runtime(..) + Self::ContractDigestParsing(..) + | Self::StateMerge(..) + | Self::StorageDecoding(..) + | Self::Runtime(..) | Self::ResultEncoding(..) | Self::CallResultEncoding(..) | Self::FrameEncoding(..) diff --git a/starstream-ledger/src/server/http/mod.rs b/starstream-ledger/src/server/http/mod.rs index 2fe6fd94..6cb99d1e 100644 --- a/starstream-ledger/src/server/http/mod.rs +++ b/starstream-ledger/src/server/http/mod.rs @@ -7,7 +7,7 @@ use core::task::{Poll, ready}; use core::time::Duration; use std::collections::{HashMap, HashSet, hash_map}; -use std::sync::{Arc, Weak}; +use std::sync::Arc; use anyhow::Context as _; use bytes::{Buf, Bytes, BytesMut}; @@ -34,17 +34,19 @@ use tokio::time::sleep; use tokio_util::codec::{Encoder as _, FramedRead}; use tokio_util::io::StreamReader; use tracing::{Instrument as _, debug, error, info, instrument, warn}; -use wasm_tokio::{AsyncReadCore as _, AsyncReadLeb128 as _, cm::U64Codec}; +use wasm_tokio::cm::{AsyncReadValue as _, U64Codec}; +use wasm_tokio::{AsyncReadCore as _, AsyncReadLeb128 as _}; use wasmparser::WasmFeatures; use wasmtime::component::{ResourceTable, Type, Val}; use wrpc_transport::FrameDecoder; +use crate::runtime::apply_state; use crate::server::{Contract, Ctx, Ledger, Transaction, UtxoCtx}; +use crate::wrpc::LEDGER_PACKAGE; use crate::wrpc::codec::{ValEncoder, read_value}; -use crate::wrpc::{LEDGER_PACKAGE, UTXO_PACKAGE}; use crate::{ APPLICATION_CBOR, APPLICATION_COSE, APPLICATION_WASM, Action, Block, Envelope, EnvelopeContext, - Fund, Publish, TransactionInput, TransactionOutput, parse_digest, + Fund, Publish, TransactionInput, parse_digest, }; mod error; @@ -476,12 +478,10 @@ impl Ledger { fn handle_genesis_get( &self, ) -> Result>, GenesisGetError> { - let genesis = - minicbor::to_vec(&self.genesis.tx_outputs).map_err(GenesisGetError::Encoding)?; http::Response::builder() .header(CONTENT_TYPE, APPLICATION_CBOR.to_string()) .header(X_CONTENT_TYPE_OPTIONS, "nosniff") - .body(http_body_util::Full::new(Bytes::from(genesis))) + .body(http_body_util::Full::new(self.genesis.encoded.clone())) .map_err(GenesisGetError::Http) } @@ -578,45 +578,28 @@ impl Ledger { } } // TODO: Verify sum(inputs) >= sum(outputs) + fee - let mut tx_outputs = Vec::with_capacity(outputs.len()); - { - let mut utxos = self.utxos.write().await; - for TransactionOutput { wasm, .. } in outputs { - let digest: [u8; 32] = Sha256::digest(&wasm).into(); - let utxo = if let Some(utxo) = utxos.get(&digest).and_then(Weak::upgrade) { - utxo - } else { - let utxo = Arc::new(wasm.into()); - utxos.insert(digest, Arc::downgrade(&utxo)); - utxo - }; - tx_outputs.push(Some(utxo)); - } - for (tx, i) in resolved_inputs { - let utxo = if let Some(tx) = tx { - let Some(Transaction { outputs, .. }) = txs.get_mut(&tx) else { - unreachable!(); - }; - outputs[i].take() - } else { - genesis[i].take() - }; - let Some(utxo) = utxo else { + for (tx, i) in resolved_inputs { + let utxo = if let Some(tx) = tx { + let Some(Transaction { outputs, .. }) = txs.get_mut(&tx) else { unreachable!(); }; - if let Some(utxo) = Arc::into_inner(utxo) { - let digest: [u8; 32] = Sha256::digest(utxo).into(); - utxos.remove(&digest); - } - } + outputs[i].take() + } else { + genesis[i].take() + }; + let Some(..) = utxo else { + unreachable!(); + }; } - txs.insert( - digest, - Transaction { - outputs: tx_outputs, - envelope: envelope.clone(), - }, - ); + let outputs = outputs + .into_iter() + .map(|utxo| Some(Arc::new(utxo))) + .collect(); + let tx = Transaction { + outputs, + envelope: envelope.clone(), + }; + txs.insert(digest, tx); let mut blocks = self.blocks.write().await; blocks.push(Block { @@ -647,18 +630,7 @@ impl Ledger { } _ => return Err(RpcPostError::FunctionNotFound { instance, name }), }, - Some((UTXO_PACKAGE, digest)) => { - let digest = parse_digest(digest).map_err(RpcPostError::UtxoDigestParsing)?; - let wasm = { - let utxos = self.utxos.read().await; - let utxo = utxos.get(&digest).ok_or(RpcPostError::UtxoNotFound)?; - let utxo = utxo.upgrade().ok_or(RpcPostError::UtxoNotFound)?; - Bytes::clone(&utxo) - }; - - // TODO: Insert traps in place of all coordination script imports - // TODO: Merge the UTXO snapshot with the contract code - + Some((LEDGER_PACKAGE, "utxo")) => { let body = FramedRead::new(body, FrameDecoder::default()).map(|frame| { let wrpc_transport::Frame { path, data } = frame?; anyhow::ensure!(path.is_empty(), "async values not supported"); @@ -666,29 +638,49 @@ impl Ledger { }); let mut body = StreamReader::new(body.map_err(std::io::Error::other)); - // TODO: Get instance from nested interface - let mut instance = String::default(); - body.read_core_name(&mut instance) + let tx = if body + .read_option_status() .await - .map_err(RpcPostError::ParameterDecoding)?; - - let n = body + .map_err(RpcPostError::ParameterDecoding)? + { + let mut transaction = String::default(); + body.read_core_name(&mut transaction) + .await + .map_err(RpcPostError::ParameterDecoding)?; + let transaction = parse_digest(&transaction) + .map_err(RpcPostError::TransactionDigestParsing)?; + Some(transaction) + } else { + None + }; + let index = body .read_u32_leb128() .await .map_err(RpcPostError::ParameterDecoding)?; - let mut methods = HashSet::default(); - for _ in 0..n { - let mut hash = [0; 4]; - for v in &mut hash { - *v = body - .read_u64_leb128() - .await - .map_err(RpcPostError::ParameterDecoding)?; - } - let [a, b, c, d] = hash; - methods.insert((a, b, c, d)); - } - let cx = Arc::new(UtxoCtx { methods }); + let index = usize::try_from(index).map_err(|_| RpcPostError::UtxoIndexOverflow)?; + + let utxo = if let Some(transaction) = tx { + let txs = self.transactions.read().await; + let Transaction { outputs, .. } = + txs.get(&transaction).ok_or(RpcPostError::UtxoNotFound)?; + outputs.get(index).cloned() + } else { + let outputs = self.genesis.outputs.read().await; + outputs.get(index).cloned() + }; + let utxo = utxo.ok_or(RpcPostError::UtxoNotFound)?; + let utxo = utxo.as_deref().ok_or(RpcPostError::UtxoNotFound)?; + let contract = + parse_digest(&utxo.contract).map_err(RpcPostError::ContractDigestParsing)?; + let wasm = { + let contracts = self.contracts.read().await; + let contract = contracts + .get(&contract) + .ok_or(RpcPostError::ContractNotFound)?; + apply_state(&contract.wasm, &utxo.state).map_err(RpcPostError::StateMerge)? + }; + + // TODO: Insert traps in place of all coordination script imports let mut imports = HashMap::default(); let contract = self @@ -696,9 +688,9 @@ impl Ledger { .await .map_err(RpcPostError::Runtime)?; - let utxo_export = contract.get_utxo(&instance).map_err(|source| { + let utxo_export = contract.get_utxo(&utxo.instance).map_err(|source| { RpcPostError::UtxoInstanceNotFound { - instance: instance.clone(), + instance: utxo.instance.clone(), source, } })?; @@ -708,19 +700,23 @@ impl Ledger { let method_export = contract .get_utxo_method(&utxo_export, &format!("[method]utxo.{name}")) .map_err(|source| RpcPostError::UtxoMethodNotFound { - instance, + instance: utxo.instance.clone(), name, source, })?; + let methods = utxo.methods.iter().copied().collect(); + let cx = Arc::new(UtxoCtx { methods }); + let mut storage = Val::Record(Vec::default()); + // TODO: use sync API read_value( - &mut body, + &mut utxo.storage.as_ref(), &mut storage, &Type::Record(storage_export.ty().clone()), ) .await - .map_err(RpcPostError::ParameterDecoding)?; + .map_err(RpcPostError::StorageDecoding)?; let param_tys = method_export.ty().params().skip(1); let mut params = vec![Val::Bool(false); param_tys.len() + 1]; diff --git a/starstream-ledger/src/server/mod.rs b/starstream-ledger/src/server/mod.rs index 58397dcc..c68e9287 100644 --- a/starstream-ledger/src/server/mod.rs +++ b/starstream-ledger/src/server/mod.rs @@ -3,11 +3,10 @@ use core::sync::atomic::AtomicU64; use std::collections::{HashMap, HashSet}; -use std::sync::{Arc, Weak}; +use std::sync::Arc; use bytes::Bytes; use ed25519_dalek::VerifyingKey; -use sha2::{Digest as _, Sha256}; use starstream_runtime_next::{ CoordinationScriptImport, UtxoImport, get_coordination_script_instance_import, utxo_imports, }; @@ -59,13 +58,13 @@ struct Contract { #[derive(Clone, Debug)] struct Transaction { - outputs: Vec>>, + outputs: Vec>>, envelope: Bytes, } struct Genesis { - outputs: RwLock>]>>, - tx_outputs: Box<[TransactionOutput]>, + outputs: RwLock>]>>, + encoded: Bytes, } /// Starstream ledger @@ -75,7 +74,6 @@ pub struct Ledger { contracts: RwLock>>, accounts: RwLock>, transactions: RwLock>, - utxos: RwLock>>, admin: AdminAccount, genesis: Genesis, network: Arc, @@ -99,30 +97,21 @@ impl Ledger { last_nonce: AtomicU64::default(), }; let genesis = genesis.into(); - let mut utxos = HashMap::with_capacity(genesis.len()); - let mut outputs = Vec::with_capacity(genesis.len()); - for TransactionOutput { wasm, .. } in &genesis { - let digest: [u8; 32] = Sha256::digest(wasm).into(); - let utxo = if let Some(utxo) = utxos.get(&digest).and_then(Weak::upgrade) { - utxo - } else { - let utxo = Arc::new(Bytes::copy_from_slice(wasm)); - utxos.insert(digest, Arc::downgrade(&utxo)); - utxo - }; - outputs.push(Some(utxo)); - } + let encoded = minicbor::to_vec(&genesis).expect("failed to encode genesis to CBOR"); + let outputs = genesis + .into_iter() + .map(|utxo| Some(Arc::new(utxo))) + .collect(); Self { engine, blocks: RwLock::default(), contracts: RwLock::default(), accounts: RwLock::default(), transactions: RwLock::default(), - utxos: RwLock::new(utxos), admin, genesis: Genesis { - outputs: RwLock::new(outputs.into()), - tx_outputs: genesis, + outputs: RwLock::new(outputs), + encoded: encoded.into(), }, network: network.into(), max_requests: max_requests as _, diff --git a/starstream-ledger/src/wrpc.rs b/starstream-ledger/src/wrpc.rs index 921d799e..682764e7 100644 --- a/starstream-ledger/src/wrpc.rs +++ b/starstream-ledger/src/wrpc.rs @@ -5,6 +5,3 @@ pub mod codec; /// The package name used for the ledger. pub const LEDGER_PACKAGE: &str = "starstream:ledger"; - -/// The package name used for UTXOs. -pub const UTXO_PACKAGE: &str = "starstream:utxo"; diff --git a/starstream-ledger/tests/common/mod.rs b/starstream-ledger/tests/common/mod.rs index 6cc695c6..044e99a7 100644 --- a/starstream-ledger/tests/common/mod.rs +++ b/starstream-ledger/tests/common/mod.rs @@ -1,11 +1,14 @@ use core::net::{Ipv6Addr, SocketAddr}; use std::collections::BTreeSet; +use std::path::Path; use std::sync::LazyLock; -use anyhow::Context as _; +use anyhow::{Context as _, anyhow, ensure}; use ed25519_dalek::SigningKey; use sha2::{Digest as _, Sha256}; +use starstream_compiler::{ModuleGraph, TypecheckOptions, typecheck_modules}; +use starstream_types::FileSystem; use tokio::net::TcpListener; #[path = "../../../starstream-runtime-next/tests/common/mod.rs"] @@ -16,6 +19,27 @@ pub static ADMIN: LazyLock = LazyLock::new(|| SigningKey::from_bytes pub const NETWORK: &str = "starstream:test"; +fn hash_methods(names: [&str; N]) -> BTreeSet<(u64, u64, u64, u64)> { + names.map(method_hash).into_iter().collect() +} + +pub fn compile_contract_file(path: impl AsRef) -> anyhow::Result> { + let mut fs = FileSystem::default(); + let (graph, module_id) = ModuleGraph::from_entry(&mut fs, path.as_ref()) + .map_err(|errors| anyhow!("failed to load module graph: {errors:?}"))?; + let graph = typecheck_modules(&graph, TypecheckOptions::default()) + .map_err(|failure| anyhow!("failed to typecheck program: {:?}", failure.errors))?; + let compile_result = starstream_to_wasm::compile_contract(&graph, module_id); + ensure!( + compile_result.errors.is_empty(), + "failed to compile program: {:#?}", + compile_result.errors + ); + compile_result + .to_component() + .map_err(|err| anyhow!("failed to componentize program: {err:?}")) +} + pub static SCORE_WASM: LazyLock> = LazyLock::new(|| { compile_contract(include_str!("../../../examples/score.star")) .unwrap() @@ -26,19 +50,46 @@ pub static SCORE_WASM_DIGEST: LazyLock<[u8; 32]> = LazyLock::new(|| Sha256::digest(&*SCORE_WASM).into()); pub static SCORE_EXAMPLE_METHODS: LazyLock> = LazyLock::new(|| { - [ + hash_methods([ "get_chips", "get_mult", "plus_chips", "plus_mult", "mult_mult", "finish", - ] - .map(method_hash) - .into_iter() - .collect() + ]) }); +pub static REQUIRE_HASH_PREIMAGE_WASM: LazyLock> = LazyLock::new(|| { + let mut sha256lib = std::process::Command::new(env!("CARGO")) + .args([ + "build", + "--release", + "--target", + "wasm32-unknown-unknown", + "-p", + "sha256lib", + ]) + .spawn() + .unwrap(); + let status = sha256lib.wait().unwrap(); + if !status.success() { + panic!("failed to compile `sha256lib`") + } + compile_contract_file(concat!( + env!("CARGO_MANIFEST_DIR"), + "/../examples/require_hash_preimage.star" + )) + .unwrap() + .into_boxed_slice() +}); + +pub static REQUIRE_HASH_PREIMAGE_WASM_DIGEST: LazyLock<[u8; 32]> = + LazyLock::new(|| Sha256::digest(&*REQUIRE_HASH_PREIMAGE_WASM).into()); + +pub static REQUIRE_HASH_PREIMAGE_CREATE_METHODS: LazyLock> = + LazyLock::new(|| hash_methods(["consume"])); + pub async fn free_tcp_addr() -> anyhow::Result { let lis = TcpListener::bind((Ipv6Addr::LOCALHOST, 0)) .await diff --git a/starstream-ledger/tests/imports.rs b/starstream-ledger/tests/imports.rs deleted file mode 100644 index 760281e9..00000000 --- a/starstream-ledger/tests/imports.rs +++ /dev/null @@ -1,107 +0,0 @@ -#![cfg(feature = "client")] - -use std::collections::HashMap; - -use anyhow::{bail, ensure}; -use bytes::Bytes; -use starstream_ledger::client::runtime::{ - Client, call_coordination_script, compile_component, new_contract, -}; -use starstream_ledger::{Transaction, TransactionInput, TransactionOutput, encode_digest}; - -pub mod common; -use common::*; - -struct ScoreClient; - -impl Client for ScoreClient { - async fn get_contract_wasm(&self, digest: [u8; 32]) -> anyhow::Result { - ensure!(digest == *SCORE_WASM_DIGEST, "unexpected contract digest"); - Ok(Bytes::from(SCORE_WASM.to_vec())) - } - - async fn get_input_utxo(&self, _input: &TransactionInput) -> anyhow::Result { - bail!("no inputs") - } -} - -#[tokio::test] -async fn imported_script_output() { - let mut config = wasmtime::Config::new(); - config.wasm_component_model_implements(true); - let engine = wasmtime::Engine::new(&config).unwrap(); - let wizer = wasmtime_wizer::Wizer::new(); - - let score_digest = encode_digest(&SCORE_WASM_DIGEST); - let wasm = wat::parse_str(format!( - r#"(component - (import "starstream:contract/scripts" (instance $scripts - (export "example" (external-id "{score_digest}") (func)) - )) - (core func $example (canon lower (func $scripts "example"))) - (core module $m - (import "" "example" (func $example)) - (func (export "run") (call $example)) - ) - (core instance $i (instantiate $m (with "" (instance (export "example" (func $example)))))) - (func (export "run") (canon lift (core func $i "run"))) -)"# - )) - .unwrap(); - - let mut imports = HashMap::default(); - let component = compile_component(&engine, &wizer, &wasm).unwrap(); - let contract = new_contract(&ScoreClient, &wizer, &component, None, &mut imports) - .await - .unwrap(); - let run = contract.get_coordination_script("run").unwrap(); - let Transaction { - inputs, outputs, .. - } = call_coordination_script( - &ScoreClient, - &wizer, - &contract, - &wasm, - &run, - &mut imports, - [], - &mut [], - ) - .await - .unwrap(); - assert_eq!(inputs, []); - - let mut score_imports = HashMap::default(); - let score_component = compile_component(&engine, &wizer, &SCORE_WASM).unwrap(); - let score = new_contract( - &ScoreClient, - &wizer, - &score_component, - None, - &mut score_imports, - ) - .await - .unwrap(); - let example = score.get_coordination_script("example").unwrap(); - let Transaction { - outputs: score_outputs, - .. - } = call_coordination_script( - &ScoreClient, - &wizer, - &score, - &SCORE_WASM, - &example, - &mut score_imports, - [], - &mut [], - ) - .await - .unwrap(); - - let [TransactionOutput { contract, .. }] = outputs.as_slice() else { - panic!("invalid outputs: {outputs:?}") - }; - assert_eq!(contract.as_ref(), score_digest); - assert_eq!(outputs, score_outputs); -} diff --git a/starstream-ledger/tests/ledger.rs b/starstream-ledger/tests/ledger.rs index 0e3df8af..9263da9d 100644 --- a/starstream-ledger/tests/ledger.rs +++ b/starstream-ledger/tests/ledger.rs @@ -25,6 +25,7 @@ use starstream_ledger::wrpc::codec::ValEncoder; use starstream_ledger::{Transaction, TransactionInput, TransactionOutput, encode_digest}; use tokio::io::AsyncReadExt as _; use tokio_util::codec::Encoder as _; +use wasmtime::Store; use wasmtime::component::{Component, Type, Val}; use wasmtime_wizer::WasmtimeWizerComponent; @@ -64,7 +65,7 @@ async fn http() { let score_example_export = score.get_coordination_script("example").unwrap(); let scope_progress_utxo_export = score.get_utxo("score-progress").unwrap(); let score_progress_utxo_storage_export = scope_progress_utxo_export.storage().unwrap(); - let mut store = wasmtime::Store::new(&engine, Ctx::default()); + let mut store = Store::new(&engine, Ctx::default()); let score = score.instantiate(&mut store).await.unwrap(); score .call_coordination_script(&mut store, &score_example_export, &[], &mut []) @@ -104,7 +105,9 @@ async fn http() { instance: "score-progress".into(), methods: SCORE_EXAMPLE_METHODS.clone(), storage: score_progress_utxo_storage_buf.to_vec().into(), - wasm: score_progress_utxo.into(), + state: starstream_ledger::runtime::parse_state(&score_progress_utxo) + .collect::>() + .unwrap(), }; let genesis = [ score_progress_genesis_utxo.clone(), @@ -313,6 +316,7 @@ async fn http() { .await .unwrap(); let score_example_export = score_contract.get_coordination_script("example").unwrap(); + let mut utxos = Vec::default(); let Transaction { inputs, outputs, @@ -320,12 +324,17 @@ async fn http() { proof, } = client .call_coordination_script( + &mut Store::new( + client.engine(), + starstream_ledger::client::runtime::Ctx::default(), + ), &score_contract, &SCORE_WASM, &score_example_export, &mut HashMap::default(), [], &mut [], + &mut utxos, ) .await .unwrap(); @@ -336,7 +345,7 @@ async fn http() { ref instance, ref methods, ref storage, - wasm: ref utxo_wasm, + .. }, ] = *outputs else { @@ -407,27 +416,32 @@ async fn http() { .expect_err("spent inputs must be rejected"); assert_eq!(err.to_string(), "input not found"); - let genesis_utxo_digest = Sha256::digest(&genesis[1].wasm).into(); let mut rx = client .call_utxo_method( - &genesis_utxo_digest, - &genesis[1].instance, + &TransactionInput { + transaction: Box::default(), + index: 1, + }, "get-chips", - &genesis[1].methods, - &genesis[1].storage, &[], ) .await - .expect("unspent genesis output must remain callable"); + .expect("genesis output must remain callable"); let mut buf = Vec::default(); rx.read_to_end(&mut buf).await.unwrap(); assert_eq!(buf, [42]); - let utxo_digest = Sha256::digest(utxo_wasm).into(); let mut rx = client - .call_utxo_method(&utxo_digest, instance, "get-chips", methods, storage, &[]) + .call_utxo_method( + &TransactionInput { + transaction: encode_digest(&tx_digest).into(), + index: 0, + }, + "get-chips", + &[], + ) .await - .expect("should succeed, because UTXO is present in genesis"); + .expect("transaction output must be callable"); let mut buf = Vec::default(); rx.read_to_end(&mut buf).await.unwrap(); assert_eq!(buf, [42]); diff --git a/starstream-runtime-next/src/lib.rs b/starstream-runtime-next/src/lib.rs index a5ab4faf..9af9ff7a 100644 --- a/starstream-runtime-next/src/lib.rs +++ b/starstream-runtime-next/src/lib.rs @@ -414,32 +414,41 @@ fn link_typed_utxo_main( fn link_typed_utxo_method( linker: &mut LinkerInstance, ty: types::ComponentFunc, - idx: ComponentExportIndex, name: &str, ) -> wasmtime::Result<()> { let Some((_, Type::Borrow(..))) = ty.params().next() else { bail!("function does not take borrowed resource type as first parameter"); }; - linker.func_new_async(name, move |mut store, _ty, params, results| { - Box::new(async move { - let Some(Val::Resource(utxo)) = params.first() else { - bail!("first parameter is not a resource") - }; - let utxo = utxo.try_into_resource::>(&mut store)?; - let &Utxo { - instance, resource, .. - } = store.data_mut().table().get(&utxo)?; - let params = { - let mut ps = Vec::with_capacity(params.len()); - ps.push(Val::Resource(resource)); - for p in ¶ms[1..] { - ps.push(p.clone()); - } - ps - }; - call_func(&mut store, &instance, idx, ¶ms, results).await?; - Ok(()) - }) + linker.func_new_async(name, { + let name = Arc::::from(name); + move |mut store, _ty, params, results| { + let name = Arc::clone(&name); + Box::new(async move { + let Some(Val::Resource(utxo)) = params.first() else { + bail!("first parameter is not a resource") + }; + let utxo = utxo.try_into_resource::>(&mut store)?; + let &Utxo { + instance, + instance_idx, + resource, + .. + } = store.data_mut().table().get(&utxo)?; + let idx = instance + .get_export_index(&mut store, Some(&instance_idx), name.as_ref()) + .with_context(|| format!("`{name}` export was not found"))?; + let params = { + let mut ps = Vec::with_capacity(params.len()); + ps.push(Val::Resource(resource)); + for p in ¶ms[1..] { + ps.push(p.clone()); + } + ps + }; + call_func(&mut store, &instance, idx, ¶ms, results).await?; + Ok(()) + }) + } }) } @@ -454,22 +463,24 @@ fn link_typed_utxo_function( name: &str, external_id: &Option>, ) -> wasmtime::Result<()> { - let idx = target - .component() - .get_export_index(Some(instance_idx), name) - .with_context(|| format!("`{name}` export was not found"))?; match name.split_once(']') { - Some(("[static", ..)) => link_typed_utxo_main( - target, - linker, - ty, - *instance_idx, - instance_name, - idx, - name, - external_id, - ), - Some(("[method", ..)) => link_typed_utxo_method(linker, ty, idx, name), + Some(("[static", ..)) => { + let idx = target + .component() + .get_export_index(Some(instance_idx), name) + .with_context(|| format!("`{name}` export was not found"))?; + link_typed_utxo_main( + target, + linker, + ty, + *instance_idx, + instance_name, + idx, + name, + external_id, + ) + } + Some(("[method", ..)) => link_typed_utxo_method(linker, ty, name), _ => bail!("unexpected typed UTXO instance function import `{name}`"), } } diff --git a/starstream-runtime-next/tests/score.rs b/starstream-runtime-next/tests/score.rs index 2fae59a9..a25ab542 100644 --- a/starstream-runtime-next/tests/score.rs +++ b/starstream-runtime-next/tests/score.rs @@ -202,6 +202,7 @@ async fn get_progress_storage( Ok(storage.iter().collect()) } +#[allow(clippy::too_many_arguments)] async fn assert_call_method( mut store: &mut Store, contract: &Contract,