diff --git a/Cargo.lock b/Cargo.lock index e7f76760a..b1a99825a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -151,8 +151,8 @@ dependencies = [ "pin-project-lite", "rustversion", "serde", - "sync_wrapper", - "tower", + "sync_wrapper 0.1.2", + "tower 0.4.13", "tower-layer", "tower-service", ] @@ -507,6 +507,23 @@ version = "1.0.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9330f8b2ff13f34540b44e946ef35111825727b38d33286ef986142615121801" +[[package]] +name = "cfg_aliases" +version = "0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f079e83a288787bcd14a6aea84cee5c87a67c5a3e660c30f557a3d24761b3527" + +[[package]] +name = "chacha20" +version = "0.10.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d524456ba66e72eb8b115ff89e01e497f8e6d11d78b70b1aa13c0fbd97540a81" +dependencies = [ + "cfg-if", + "cpufeatures", + "rand_core 0.10.1", +] + [[package]] name = "chrono" version = "0.4.44" @@ -726,6 +743,15 @@ dependencies = [ "serde_json", ] +[[package]] +name = "cpufeatures" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8b2a41393f66f16b0823bb79094d54ac5fbd34ab292ddafb9a0456ac9f87d201" +dependencies = [ + "libc", +] + [[package]] name = "crc32fast" version = "1.5.0" @@ -870,7 +896,7 @@ dependencies = [ "libc", "option-ext", "redox_users", - "windows-sys 0.59.0", + "windows-sys 0.61.2", ] [[package]] @@ -942,6 +968,7 @@ dependencies = [ "prometheus", "rand 0.9.2", "rayon", + "reqwest 0.12.28", "rocksdb", "serde", "serde_derive", @@ -1060,7 +1087,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.59.0", + "windows-sys 0.61.2", ] [[package]] @@ -1215,8 +1242,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ff2abc00be7fca6ebc474524697ae276ad847ad0a6b3faa4bcb027e9a4614ad0" dependencies = [ "cfg-if", + "js-sys", "libc", "wasi", + "wasm-bindgen", ] [[package]] @@ -1238,10 +1267,13 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "139ef39800118c7683f2fd3c98c1b23c09ae076556b435f8e9064ae108aaeeec" dependencies = [ "cfg-if", + "js-sys", "libc", "r-efi", + "rand_core 0.10.1", "wasip2", "wasip3", + "wasm-bindgen", ] [[package]] @@ -1462,6 +1494,23 @@ dependencies = [ "pin-utils", "smallvec", "tokio", + "want", +] + +[[package]] +name = "hyper-rustls" +version = "0.27.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "33ca68d021ef39cf6463ab54c1d0f5daf03377b70561305bb89a8f83aab66e0f" +dependencies = [ + "http 1.4.0", + "hyper 1.8.1", + "hyper-util", + "rustls 0.23.43", + "tokio", + "tokio-rustls", + "tower-service", + "webpki-roots 1.0.9", ] [[package]] @@ -1482,12 +1531,21 @@ version = "0.1.20" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "96547c2556ec9d12fb1578c4eaf448b04993e7fb79cbaad930a656880a6bdfa0" dependencies = [ + "base64 0.22.1", "bytes", + "futures-channel", + "futures-util", "http 1.4.0", "http-body 1.0.1", "hyper 1.8.1", + "ipnet", + "libc", + "percent-encoding", "pin-project-lite", + "socket2 0.6.3", "tokio", + "tower-service", + "tracing", ] [[package]] @@ -1669,7 +1727,7 @@ checksum = "3640c1c38b8e4e43584d8df18be5fc6b0aa314ce6ebf51b53313d4306cca8e46" dependencies = [ "hermit-abi 0.5.2", "libc", - "windows-sys 0.59.0", + "windows-sys 0.61.2", ] [[package]] @@ -1873,6 +1931,12 @@ version = "0.4.29" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5e5032e24019045c762d3c0f28f5b6b8bbf38563a65908389bf7978758920897" +[[package]] +name = "lru-slab" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "112b39cec0b298b6c1999fee3e31427f74f676e4cb9879ed1a121b43661a4154" + [[package]] name = "lz4-sys" version = "1.11.1+lz4-1.10.0" @@ -1939,7 +2003,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "05015102dad0f7d61691ca347e9d9d9006685a64aefb3d79eecf62665de2153d" dependencies = [ "rustls 0.21.12", - "rustls-webpki", + "rustls-webpki 0.101.7", "serde", "serde_json", "webpki-roots 0.25.4", @@ -2032,7 +2096,7 @@ dependencies = [ "bytes", "http 0.2.12", "opentelemetry_api", - "reqwest", + "reqwest 0.11.27", ] [[package]] @@ -2050,7 +2114,7 @@ dependencies = [ "opentelemetry_api", "opentelemetry_sdk", "prost", - "reqwest", + "reqwest 0.11.27", "thiserror 1.0.69", "tokio", "tonic", @@ -2368,6 +2432,62 @@ dependencies = [ "percent-encoding", ] +[[package]] +name = "quinn" +version = "0.11.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0c1a41e437b6bbd489372cd4971de128e85c855f56c57f283d20ff016cf7c0a8" +dependencies = [ + "bytes", + "cfg_aliases", + "pin-project-lite", + "quinn-proto", + "quinn-udp", + "rustc-hash", + "rustls 0.23.43", + "socket2 0.6.3", + "thiserror 2.0.18", + "tokio", + "tracing", + "web-time", +] + +[[package]] +name = "quinn-proto" +version = "0.11.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "04759210543be93709136e28212294a659ef5001836ff4eab4d663e4529bba83" +dependencies = [ + "bytes", + "getrandom 0.4.1", + "lru-slab", + "rand 0.10.2", + "rand_pcg", + "ring 0.17.14", + "rustc-hash", + "rustls 0.23.43", + "rustls-pki-types", + "slab", + "thiserror 2.0.18", + "tinyvec", + "tracing", + "web-time", +] + +[[package]] +name = "quinn-udp" +version = "0.5.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "35a133f956daabe89a61a685c2649f13d82d5aa4bd5d12d1277e1072a21c0694" +dependencies = [ + "cfg_aliases", + "libc", + "once_cell", + "socket2 0.6.3", + "tracing", + "windows-sys 0.61.2", +] + [[package]] name = "quote" version = "1.0.44" @@ -2404,6 +2524,17 @@ dependencies = [ "rand_core 0.9.5", ] +[[package]] +name = "rand" +version = "0.10.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c7f5fa3a058cd35567ef9bfa5e75732bee0f9e4c55fa90477bef2dfcdbc4be80" +dependencies = [ + "chacha20", + "getrandom 0.4.1", + "rand_core 0.10.1", +] + [[package]] name = "rand_chacha" version = "0.3.1" @@ -2442,6 +2573,21 @@ dependencies = [ "getrandom 0.3.4", ] +[[package]] +name = "rand_core" +version = "0.10.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "63b8176103e19a2643978565ca18b50549f6101881c443590420e4dc998a3c69" + +[[package]] +name = "rand_pcg" +version = "0.10.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "caa0f4137e1c0a72f4c651489402276c8e8e1cf081f3b0ba156d2cbeef09e86a" +dependencies = [ + "rand_core 0.10.1", +] + [[package]] name = "rayon" version = "1.11.0" @@ -2545,7 +2691,7 @@ dependencies = [ "serde", "serde_json", "serde_urlencoded", - "sync_wrapper", + "sync_wrapper 0.1.2", "system-configuration", "tokio", "tower-service", @@ -2556,6 +2702,44 @@ dependencies = [ "winreg", ] +[[package]] +name = "reqwest" +version = "0.12.28" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "eddd3ca559203180a307f12d114c268abf583f59b03cb906fd0b3ff8646c1147" +dependencies = [ + "base64 0.22.1", + "bytes", + "futures-core", + "http 1.4.0", + "http-body 1.0.1", + "http-body-util", + "hyper 1.8.1", + "hyper-rustls", + "hyper-util", + "js-sys", + "log", + "percent-encoding", + "pin-project-lite", + "quinn", + "rustls 0.23.43", + "rustls-pki-types", + "serde", + "serde_json", + "serde_urlencoded", + "sync_wrapper 1.0.2", + "tokio", + "tokio-rustls", + "tower 0.5.3", + "tower-http", + "tower-service", + "url", + "wasm-bindgen", + "wasm-bindgen-futures", + "web-sys", + "webpki-roots 1.0.9", +] + [[package]] name = "ring" version = "0.16.20" @@ -2639,7 +2823,7 @@ dependencies = [ "errno 0.3.14", "libc", "linux-raw-sys 0.12.1", - "windows-sys 0.59.0", + "windows-sys 0.61.2", ] [[package]] @@ -2676,10 +2860,34 @@ checksum = "3f56a14d1f48b391359b22f731fd4bd7e43c97f3c50eee276f3aa09c94784d3e" dependencies = [ "log", "ring 0.17.14", - "rustls-webpki", + "rustls-webpki 0.101.7", "sct 0.7.1", ] +[[package]] +name = "rustls" +version = "0.23.43" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0283386ce02abc0151e1761d08802dfe86c173b0b494af5cbc086574e453da06" +dependencies = [ + "once_cell", + "ring 0.17.14", + "rustls-pki-types", + "rustls-webpki 0.103.14", + "subtle", + "zeroize", +] + +[[package]] +name = "rustls-pki-types" +version = "1.15.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2f4925028c7eb5d1fcdaf196971378ed9d2c1c4efc7dc5d011256f76c99c0a96" +dependencies = [ + "web-time", + "zeroize", +] + [[package]] name = "rustls-webpki" version = "0.101.7" @@ -2690,6 +2898,17 @@ dependencies = [ "untrusted 0.9.0", ] +[[package]] +name = "rustls-webpki" +version = "0.103.14" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0527518605e68109d875e248ea259b6758801cf165e4b2c2733ae3b51f12535a" +dependencies = [ + "ring 0.17.14", + "rustls-pki-types", + "untrusted 0.9.0", +] + [[package]] name = "rustversion" version = "1.0.22" @@ -2931,7 +3150,7 @@ version = "1.4.8" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c4db69cba1110affc0e9f7bcd48bbf87b3f4fc7c61fc9155afd4c469eb3d6c1b" dependencies = [ - "errno 0.2.8", + "errno 0.3.14", "libc", ] @@ -3073,6 +3292,12 @@ version = "0.8.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8ea5119cdb4c55b55d432abb513a0429384878c15dde60cc77b1c99de1a95a6a" +[[package]] +name = "subtle" +version = "2.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "13c2bddecc57b384dee18652358fb23172facb8a2c51ccc10d74c157bdea3292" + [[package]] name = "syn" version = "1.0.109" @@ -3101,6 +3326,15 @@ version = "0.1.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2047c6ded9c721764247e62cd3b03c09ffc529b2ba5b10ec482ae507a4a70160" +[[package]] +name = "sync_wrapper" +version = "1.0.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0bf256ce5efdfa370213c1dabab5935a12e49f2c58d15e9eac2870d3b4f27263" +dependencies = [ + "futures-core", +] + [[package]] name = "synstructure" version = "0.13.2" @@ -3185,7 +3419,7 @@ dependencies = [ "getrandom 0.4.1", "once_cell", "rustix 1.1.4", - "windows-sys 0.59.0", + "windows-sys 0.61.2", ] [[package]] @@ -3427,6 +3661,16 @@ dependencies = [ "syn 2.0.117", ] +[[package]] +name = "tokio-rustls" +version = "0.26.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1729aa945f29d91ba541258c8df89027d5792d85a8841fb65e8bf0f4ede4ef61" +dependencies = [ + "rustls 0.23.43", + "tokio", +] + [[package]] name = "tokio-stream" version = "0.1.18" @@ -3507,7 +3751,7 @@ dependencies = [ "prost", "tokio", "tokio-stream", - "tower", + "tower 0.4.13", "tower-layer", "tower-service", "tracing", @@ -3533,6 +3777,39 @@ dependencies = [ "tracing", ] +[[package]] +name = "tower" +version = "0.5.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ebe5ef63511595f1344e2d5cfa636d973292adc0eec1f0ad45fae9f0851ab1d4" +dependencies = [ + "futures-core", + "futures-util", + "pin-project-lite", + "sync_wrapper 1.0.2", + "tokio", + "tower-layer", + "tower-service", +] + +[[package]] +name = "tower-http" +version = "0.6.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4cfcf7e2740e6fc6d4d688b4ef00650406bb94adf4731e43c096c3a19fe40840" +dependencies = [ + "bitflags 2.11.0", + "bytes", + "futures-util", + "http 1.4.0", + "http-body 1.0.1", + "pin-project-lite", + "tower 0.5.3", + "tower-layer", + "tower-service", + "url", +] + [[package]] name = "tower-layer" version = "0.3.3" @@ -3923,6 +4200,16 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "web-time" +version = "1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5a6580f308b1fad9207618087a65c04e7a10bc77e02c8e84e9b00dd4b12fa0bb" +dependencies = [ + "js-sys", + "wasm-bindgen", +] + [[package]] name = "webpki" version = "0.21.4" @@ -3957,6 +4244,15 @@ version = "0.25.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5f20c57d8d7db6d3b86154206ae5d8fba62dd39573114de97c2cb0578251f8e1" +[[package]] +name = "webpki-roots" +version = "1.0.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7dcd9d09a39985f5344844e66b0c530a33843579125f23e21e9f0f220850f22a" +dependencies = [ + "rustls-pki-types", +] + [[package]] name = "which" version = "3.1.1" @@ -4012,7 +4308,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.48.0", + "windows-sys 0.61.2", ] [[package]] @@ -4424,6 +4720,12 @@ dependencies = [ "synstructure", ] +[[package]] +name = "zeroize" +version = "1.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e13c156562582aa81c60cb29407084cdb54c4164760106ab78e6c5b0858cf64e" + [[package]] name = "zeromq-src" version = "0.2.6+4.3.4" diff --git a/Cargo.toml b/Cargo.toml index 828ebe165..60f74133a 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -17,7 +17,7 @@ default-run = "electrs" unexpected_cfgs = { level = "warn", check-cfg = ["cfg(has_error_description_deprecated)"] } [features] -liquid = ["elements"] +liquid = ["elements", "reqwest"] electrum-discovery = ["electrum-client"] bench = [] otlp-tracing = [ @@ -59,7 +59,7 @@ serde_json = "1.0.60" signal-hook = "0.4" stderrlog = "0.6" sysconf = ">=0.3.4" -time = { version = "0.3", features = ["formatting"] } +time = { version = "0.3", features = ["formatting", "parsing"] } tiny_http = "0.12.0" url = "2.2.0" hyper = { version = "1", features = ["http1", "server"] } @@ -73,6 +73,7 @@ tracing-subscriber = { version = "0.3.17", default-features = false, features = opentelemetry-semantic-conventions = { version = "0.12.0", optional = true } tracing = { version = "0.1.40", default-features = false, features = ["attributes"], optional = true } rand = "0.9.1" +reqwest = { version = "0.12", default-features = false, features = ["json", "rustls-tls"], optional = true } # optional dependencies for electrum-discovery electrum-client = { version = "0.8", optional = true } diff --git a/README.md b/README.md index 5e911aa16..a4c243784 100644 --- a/README.md +++ b/README.md @@ -74,6 +74,12 @@ In addition to electrs's original configuration options, a few new options are a Additional options with the `liquid` feature: - `--parent-network ` - the parent network this chain is pegged to. +- `--asset-registry-url ` - base URL for the Liquid asset registry v2 service. + +Asset registry timeouts can be configured in milliseconds with the +`ELECTRS_ASSET_REGISTRY_CONNECT_TIMEOUT_MS` and +`ELECTRS_ASSET_REGISTRY_REQUEST_TIMEOUT_MS` environment variables. They default +to 2000 ms and 5000 ms, respectively. Additional options with the `electrum-discovery` feature: - `--electrum-hosts ` - a json map of the public hosts where the electrum server is reachable, in the [`server.features` format](https://electrum-protocol.readthedocs.io/en/latest/protocol-methods.html#server-features). diff --git a/src/bin/electrs.rs b/src/bin/electrs.rs index 74d06f806..3ea4e993a 100644 --- a/src/bin/electrs.rs +++ b/src/bin/electrs.rs @@ -29,7 +29,7 @@ use electrs::{ use electrs::otlp_trace; #[cfg(feature = "liquid")] -use electrs::elements::AssetRegistry; +use electrs::elements::RegistryClient; use electrs::metrics::MetricOpts; /// Default salt rotation interval in seconds (24 hours) @@ -115,11 +115,13 @@ fn run_server(config: Arc, salt_rwlock: Arc>) -> Result<( } #[cfg(feature = "liquid")] - let asset_db = config.asset_db_path.as_ref().map(|db_dir| { - let asset_db = Arc::new(RwLock::new(AssetRegistry::new(db_dir.clone()))); - AssetRegistry::spawn_sync(asset_db.clone()); - asset_db - }); + let asset_registry = config + .asset_registry_url + .clone() + .map(RegistryClient::new) + .transpose() + .chain_err(|| "failed creating asset registry client")? + .map(Arc::new); let query = Arc::new(Query::new( Arc::clone(&chain), @@ -127,7 +129,7 @@ fn run_server(config: Arc, salt_rwlock: Arc>) -> Result<( Arc::clone(&daemon), Arc::clone(&config), #[cfg(feature = "liquid")] - asset_db, + asset_registry, )); // TODO: configuration for which servers to start diff --git a/src/config.rs b/src/config.rs index c1b94a73d..eb0bd58a5 100644 --- a/src/config.rs +++ b/src/config.rs @@ -7,6 +7,8 @@ use std::path::{Path, PathBuf}; use std::sync::Arc; use std::time::Duration; use stderrlog; +#[cfg(feature = "liquid")] +use url::Url; use crate::chain::Network; use crate::daemon::CookieGetter; @@ -86,7 +88,7 @@ pub struct Config { #[cfg(feature = "liquid")] pub parent_network: BNetwork, #[cfg(feature = "liquid")] - pub asset_db_path: Option, + pub asset_registry_url: Option, #[cfg(feature = "electrum-discovery")] pub electrum_public_hosts: Option, @@ -328,9 +330,9 @@ impl Config { .takes_value(true), ) .arg( - Arg::with_name("asset_db_path") - .long("asset-db-path") - .help("Directory for liquid/elements asset db") + Arg::with_name("asset_registry_url") + .long("asset-registry-url") + .help("Base URL for the Liquid asset registry v2 service") .takes_value(true), ); @@ -369,7 +371,14 @@ impl Config { }); #[cfg(feature = "liquid")] - let asset_db_path = m.value_of("asset_db_path").map(PathBuf::from); + let asset_registry_url = m.value_of("asset_registry_url").map(|value| { + let url = Url::parse(value).expect("invalid asset registry URL"); + assert!( + matches!(url.scheme(), "http" | "https"), + "asset registry URL must use http or https" + ); + url + }); let default_daemon_port = match network_type { #[cfg(not(feature = "liquid"))] @@ -572,7 +581,7 @@ impl Config { #[cfg(feature = "liquid")] parent_network, #[cfg(feature = "liquid")] - asset_db_path, + asset_registry_url, #[cfg(feature = "electrum-discovery")] electrum_public_hosts, diff --git a/src/elements/asset.rs b/src/elements/asset.rs index d47b9bc4c..9c0b7ffff 100644 --- a/src/elements/asset.rs +++ b/src/elements/asset.rs @@ -1,5 +1,4 @@ use std::collections::{HashMap, HashSet}; -use std::sync::{Arc, RwLock}; use bitcoin::hashes::{sha256, Hash}; use elements::confidential::{Asset, Value}; @@ -9,7 +8,7 @@ use elements::{issuance::ContractHash, AssetId, AssetIssuance, OutPoint, Transac use crate::chain::{BNetwork, BlockHash, Network, Txid}; use crate::elements::peg::{get_pegin_data, get_pegout_data, PeginInfo, PegoutInfo}; -use crate::elements::registry::{AssetMeta, AssetRegistry}; +use crate::elements::registry::AssetMeta; use crate::errors::*; use crate::new_index::schema::{TxHistoryInfo, TxHistoryKey, TxHistoryRow}; use crate::new_index::{db::DBFlush, ChainQuery, DBRow, Mempool, Query}; @@ -351,9 +350,8 @@ fn asset_history_row( pub fn lookup_asset( query: &Query, - registry: Option<&Arc>>, asset_id: &AssetId, - meta: Option<&AssetMeta>, // may optionally be provided if already known + meta: Option, ) -> Result> { if query.network().pegged_asset() == Some(asset_id) { let (chain_stats, mempool_stats) = pegged_asset_stats(query, asset_id); @@ -380,9 +378,6 @@ pub fn lookup_asset( Ok(if let Some(row) = row { let reissuance_token = parse_asset_id(&row.reissuance_token); - let meta = meta - .cloned() - .or_else(|| registry.and_then(|r| r.read().unwrap().get(asset_id).cloned())); let stats = issued_asset_stats(query.chain(), &mempool, asset_id, &reissuance_token); let status = query.get_tx_status(&deserialize(&row.issuance_txid).unwrap()); diff --git a/src/elements/mod.rs b/src/elements/mod.rs index e0d044c2b..d7efd845d 100644 --- a/src/elements/mod.rs +++ b/src/elements/mod.rs @@ -8,7 +8,10 @@ mod registry; use asset::get_issuance_entropy; pub use asset::{lookup_asset, LiquidAsset}; -pub use registry::{AssetRegistry, AssetSorting}; +pub use registry::{ + AssetMeta, AssetSearchFilters, AssetSorting, RegistryAsset, RegistryAssetList, RegistryClient, + RegistryContract, RegistryError, RegistryIcon, +}; #[derive(Serialize, Deserialize, Clone)] pub struct IssuanceValue { diff --git a/src/elements/registry.rs b/src/elements/registry.rs index 3aca4bff3..62498d9ae 100644 --- a/src/elements/registry.rs +++ b/src/elements/registry.rs @@ -1,123 +1,643 @@ use std::collections::HashMap; -use std::str::FromStr; -use std::sync::{Arc, RwLock}; -use std::time::{Duration, SystemTime}; -use std::{cmp, fs, path, thread}; - -use serde_json::Value as JsonValue; +use std::env; +use std::fmt; +use std::sync::Arc; +use std::time::{Duration, Instant}; use elements::AssetId; +use reqwest::{Client, StatusCode}; +use serde::de::DeserializeOwned; +use serde_json::{Map as JsonMap, Value as JsonValue}; +use time::{format_description::well_known::Rfc3339, OffsetDateTime}; +use tokio::sync::{oneshot, Mutex, OwnedSemaphorePermit, Semaphore}; +use url::Url; use crate::errors::*; -// length of asset id prefix to use for sub-directory partitioning -// (in number of hex characters, not bytes) +const DEFAULT_REGISTRY_CONNECT_TIMEOUT: Duration = Duration::from_secs(2); +const DEFAULT_REGISTRY_REQUEST_TIMEOUT: Duration = Duration::from_secs(5); +const REGISTRY_CONNECT_TIMEOUT_ENV: &str = "ELECTRS_ASSET_REGISTRY_CONNECT_TIMEOUT_MS"; +const REGISTRY_REQUEST_TIMEOUT_ENV: &str = "ELECTRS_ASSET_REGISTRY_REQUEST_TIMEOUT_MS"; +const REGISTRY_MAX_PAGE_SIZE: usize = 500; +const REGISTRY_MAX_PAGE: usize = 1_000_000; +const REGISTRY_ASSET_CACHE_TTL: Duration = Duration::from_secs(1); +const REGISTRY_ASSET_CACHE_MAX_ENTRIES: usize = 1024; +const REGISTRY_MAX_CONCURRENT_REQUESTS: usize = 16; +const REGISTRY_MAX_ASSET_RESPONSE_SIZE: usize = 1024 * 1024; +const REGISTRY_MAX_LIST_RESPONSE_SIZE: usize = 16 * 1024 * 1024; + +#[derive(Clone, Debug)] +pub enum RegistryError { + InvalidBaseUrl(String), + InvalidRequest(String), + Timeout(String), + Transport(String), + HttpStatus(u16), + InvalidResponse(String), + Overloaded(String), + LocalLookup(String), +} + +impl fmt::Display for RegistryError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::InvalidBaseUrl(message) + | Self::InvalidRequest(message) + | Self::Timeout(message) + | Self::Transport(message) + | Self::InvalidResponse(message) + | Self::Overloaded(message) + | Self::LocalLookup(message) => f.write_str(message), + Self::HttpStatus(status) => { + write!(f, "asset registry returned HTTP status {}", status) + } + } + } +} + +impl std::error::Error for RegistryError {} + +type RegistryAssetResult = std::result::Result, RegistryError>; + +enum AssetCacheEntry { + Ready { + fetched_at: Instant, + asset: Option>, + }, + Fetching(Vec>), +} -const DIR_PARTITION_LEN: usize = 2; -pub struct AssetRegistry { - directory: path::PathBuf, - assets_cache: HashMap, +#[derive(Clone)] +pub struct RegistryClient { + base_url: Url, + http: Client, + asset_cache: Arc>>, + concurrency: Arc, + asset_cache_ttl: Duration, + asset_cache_max_entries: usize, } -pub type AssetEntry<'a> = (&'a AssetId, &'a AssetMeta); +impl RegistryClient { + pub fn new(base_url: Url) -> std::result::Result { + Self::with_timeouts( + base_url, + registry_timeout_from_env( + REGISTRY_CONNECT_TIMEOUT_ENV, + DEFAULT_REGISTRY_CONNECT_TIMEOUT, + )?, + registry_timeout_from_env( + REGISTRY_REQUEST_TIMEOUT_ENV, + DEFAULT_REGISTRY_REQUEST_TIMEOUT, + )?, + ) + } + + fn with_timeouts( + base_url: Url, + connect_timeout: Duration, + request_timeout: Duration, + ) -> std::result::Result { + Self::with_options( + base_url, + connect_timeout, + request_timeout, + REGISTRY_ASSET_CACHE_TTL, + REGISTRY_ASSET_CACHE_MAX_ENTRIES, + REGISTRY_MAX_CONCURRENT_REQUESTS, + ) + } + + fn with_options( + mut base_url: Url, + connect_timeout: Duration, + request_timeout: Duration, + asset_cache_ttl: Duration, + asset_cache_max_entries: usize, + max_concurrent_requests: usize, + ) -> std::result::Result { + if !matches!(base_url.scheme(), "http" | "https") { + return Err(RegistryError::InvalidBaseUrl(format!( + "asset registry URL must use http or https: {}", + base_url + ))); + } -impl AssetRegistry { - pub fn new(directory: path::PathBuf) -> Self { - Self { - directory, - assets_cache: Default::default(), + if !base_url.path().ends_with('/') { + let path = format!("{}/", base_url.path()); + base_url.set_path(&path); } + base_url.set_query(None); + base_url.set_fragment(None); + + let http = Client::builder() + .connect_timeout(connect_timeout) + .timeout(request_timeout) + .redirect(reqwest::redirect::Policy::none()) + .user_agent(concat!("electrs/", env!("CARGO_PKG_VERSION"))) + .build() + .map_err(|error| RegistryError::Transport(error.to_string()))?; + + Ok(Self { + base_url, + http, + asset_cache: Arc::new(Mutex::new(HashMap::new())), + concurrency: Arc::new(Semaphore::new(max_concurrent_requests.max(1))), + asset_cache_ttl, + asset_cache_max_entries: asset_cache_max_entries.max(1), + }) } - pub fn get(&self, asset_id: &AssetId) -> Option<&AssetMeta> { - self.assets_cache - .get(asset_id) - .map(|(_, metadata)| metadata) + pub async fn get_asset( + &self, + asset_id: &AssetId, + ) -> std::result::Result, RegistryError> { + let mut cache = self.asset_cache.lock().await; + let now = Instant::now(); + + if let Some(AssetCacheEntry::Ready { fetched_at, asset }) = cache.get(asset_id) { + if now.duration_since(*fetched_at) < self.asset_cache_ttl { + return Ok(asset.as_ref().map(|asset| asset.as_ref().clone())); + } + } + + let receiver = if let Some(AssetCacheEntry::Fetching(waiters)) = cache.get_mut(asset_id) { + let (sender, receiver) = oneshot::channel(); + waiters.push(sender); + receiver + } else { + cache.remove(asset_id); + prune_asset_cache( + &mut cache, + now, + self.asset_cache_ttl, + self.asset_cache_max_entries, + ); + + // Same-key callers join the in-flight request above without consuming a permit. + // New requests are rejected instead of accumulating behind the registry. + let permit = self.try_acquire_permit()?; + let (sender, receiver) = oneshot::channel(); + cache.insert(*asset_id, AssetCacheEntry::Fetching(vec![sender])); + + let client = self.clone(); + let asset_id = *asset_id; + tokio::spawn(async move { + let result = client.fetch_asset(&asset_id).await; + if let Err(error) = &result { + warn!("asset registry lookup failed for {}: {}", asset_id, error); + } + client.complete_asset_fetch(asset_id, result).await; + drop(permit); + }); + receiver + }; + drop(cache); + + receiver.await.unwrap_or_else(|_| { + Err(RegistryError::Transport( + "asset registry lookup worker stopped unexpectedly".to_string(), + )) + }) + } + + async fn fetch_asset(&self, asset_id: &AssetId) -> RegistryAssetResult { + let url = self + .base_url + .join(&format!("v2/assets/{}", asset_id)) + .map_err(|error| RegistryError::InvalidBaseUrl(error.to_string()))?; + let response = self.http.get(url).send().await.map_err(request_error)?; + + if response.status() == StatusCode::NOT_FOUND { + return Ok(None); + } + if !response.status().is_success() { + return Err(RegistryError::HttpStatus(response.status().as_u16())); + } + + let mut asset: RegistryAsset = + decode_json_response(response, REGISTRY_MAX_ASSET_RESPONSE_SIZE).await?; + if asset.asset_id != *asset_id { + return Err(RegistryError::InvalidResponse(format!( + "asset registry returned {} for requested asset {}", + asset.asset_id, asset_id + ))); + } + self.make_icon_absolute(&mut asset)?; + Ok(Some(asset)) } - pub fn list( + async fn complete_asset_fetch(&self, asset_id: AssetId, result: RegistryAssetResult) { + let mut cache = self.asset_cache.lock().await; + let waiters = match cache.remove(&asset_id) { + Some(AssetCacheEntry::Fetching(waiters)) => waiters, + _ => vec![], + }; + if let Ok(asset) = &result { + cache.insert( + asset_id, + AssetCacheEntry::Ready { + fetched_at: Instant::now(), + asset: asset.clone().map(Arc::new), + }, + ); + } + drop(cache); + + for waiter in waiters { + let _ = waiter.send(result.clone()); + } + } + + fn make_icon_absolute( + &self, + asset: &mut RegistryAsset, + ) -> std::result::Result<(), RegistryError> { + if let Some(icon) = &mut asset.icon { + let href = self + .base_url + .join(&icon.href) + .map_err(|error| RegistryError::InvalidResponse(error.to_string()))?; + if !matches!(href.scheme(), "http" | "https") || href.origin() != self.base_url.origin() + { + return Err(RegistryError::InvalidResponse(format!( + "asset registry returned an invalid icon URL: {}", + icon.href + ))); + } + icon.href = href.to_string(); + } + Ok(()) + } + + fn try_acquire_permit(&self) -> std::result::Result { + self.concurrency.clone().try_acquire_owned().map_err(|_| { + RegistryError::Overloaded("too many concurrent asset registry requests".to_string()) + }) + } + + pub async fn list_assets( &self, start_index: usize, limit: usize, sorting: AssetSorting, - ) -> (usize, Vec>) { - let mut assets: Vec = self - .assets_cache - .iter() - .map(|(asset_id, (_, metadata))| (asset_id, metadata)) - .collect(); - assets.sort_by(sorting.as_comparator()); - ( - assets.len(), - assets.into_iter().skip(start_index).take(limit).collect(), - ) - } + filters: &AssetSearchFilters, + ) -> std::result::Result { + if limit > REGISTRY_MAX_PAGE_SIZE { + return Err(RegistryError::InvalidRequest(format!( + "asset registry page size cannot exceed {}", + REGISTRY_MAX_PAGE_SIZE + ))); + } + + // The v2 API has no zero-sized page, but electrs historically accepts limit=0 and + // still returns the total count. + let page_size = limit.max(1); + let page = if limit == 0 { + 1 + } else { + (start_index / page_size).checked_add(1).ok_or_else(|| { + RegistryError::InvalidRequest("asset registry page overflow".to_string()) + })? + }; + if page > REGISTRY_MAX_PAGE { + return Err(RegistryError::InvalidRequest(format!( + "asset registry page cannot exceed {}", + REGISTRY_MAX_PAGE + ))); + } + + let offset = if limit == 0 { + 0 + } else { + start_index % page_size + }; + let _permit = self.try_acquire_permit()?; + let first = self.fetch_page(page, page_size, sorting, filters).await?; + let total_count = first.total_count.ok_or_else(|| { + RegistryError::InvalidResponse( + "asset registry response is missing total_count".to_string(), + ) + })?; - pub fn fs_sync(&mut self) -> Result<()> { - let mut assets_cache = HashMap::new(); + if limit == 0 || start_index >= total_count { + return Ok(RegistryAssetList { + total_count, + items: vec![], + }); + } - for entry in fs::read_dir(&self.directory).chain_err(|| "failed reading asset dir")? { - let entry = entry.chain_err(|| "invalid fh")?; - let filetype = entry.file_type().chain_err(|| "failed getting file type")?; - if !filetype.is_dir() || entry.file_name().len() != DIR_PARTITION_LEN { - continue; + let mut items = first.items; + if offset.saturating_add(limit) > items.len() + && page < REGISTRY_MAX_PAGE + && start_index.saturating_add(items.len().saturating_sub(offset)) < total_count + { + let second = self + .fetch_page(page + 1, page_size, sorting, filters) + .await?; + if second.total_count.is_none() { + return Err(RegistryError::InvalidResponse( + "asset registry response is missing total_count".to_string(), + )); } + items.extend(second.items); + } - for file_entry in - fs::read_dir(entry.path()).chain_err(|| "failed reading asset subdir")? - { - let file_entry = file_entry.chain_err(|| "invalid fh")?; - let path = file_entry.path(); - if path.extension().and_then(|e| e.to_str()) != Some("json") { - continue; - } + Ok(RegistryAssetList { + total_count, + items: items.into_iter().skip(offset).take(limit).collect(), + }) + } - let asset_id = AssetId::from_str( - path.file_stem() - .unwrap() // cannot fail if extension() succeeded - .to_str() - .chain_err(|| "invalid filename")?, - ) - .chain_err(|| "invalid filename")?; - - let modified = file_entry - .metadata() - .chain_err(|| "failed reading metadata")? - .modified() - .chain_err(|| "metadata modified failed")?; - - if let Some((last_update, metadata)) = self.assets_cache.get(&asset_id) { - if *last_update == modified { - assets_cache.insert(asset_id, (modified, metadata.clone())); - continue; - } - } + async fn fetch_page( + &self, + page: usize, + page_size: usize, + sorting: AssetSorting, + filters: &AssetSearchFilters, + ) -> std::result::Result { + let url = self + .base_url + .join("v2/assets") + .map_err(|error| RegistryError::InvalidBaseUrl(error.to_string()))?; + let mut query = vec![ + ("page", page.to_string()), + ("page_size", page_size.to_string()), + ("sort", sorting.as_str().to_string()), + ]; + filters.append_query(&mut query); + let response = self + .http + .get(url) + .query(&query) + .send() + .await + .map_err(request_error)?; - let metadata: AssetMeta = serde_json::from_str( - &fs::read_to_string(path).chain_err(|| "failed reading file")?, - ) - .chain_err(|| "failed parsing file")?; + if !response.status().is_success() { + return Err(RegistryError::HttpStatus(response.status().as_u16())); + } - assets_cache.insert(asset_id, (modified, metadata)); - } + let mut page_response: RegistryListResponse = + decode_json_response(response, REGISTRY_MAX_LIST_RESPONSE_SIZE).await?; + if page_response.page != page || page_response.page_size != page_size { + return Err(RegistryError::InvalidResponse(format!( + "asset registry returned page {}/{} for requested page {}/{}", + page_response.page, page_response.page_size, page, page_size + ))); + } + for asset in &mut page_response.items { + self.make_icon_absolute(asset)?; } + Ok(page_response) + } +} - self.assets_cache = assets_cache; - Ok(()) +fn registry_timeout_from_env( + variable: &str, + default: Duration, +) -> std::result::Result { + match env::var(variable) { + Ok(value) => parse_registry_timeout(variable, &value), + Err(env::VarError::NotPresent) => Ok(default), + Err(env::VarError::NotUnicode(_)) => Err(RegistryError::InvalidRequest(format!( + "{} must be valid Unicode", + variable + ))), + } +} + +fn parse_registry_timeout( + variable: &str, + value: &str, +) -> std::result::Result { + let millis = value.parse::().map_err(|_| { + RegistryError::InvalidRequest(format!("{} must be a positive integer", variable)) + })?; + if millis == 0 { + return Err(RegistryError::InvalidRequest(format!( + "{} must be greater than zero", + variable + ))); } + Ok(Duration::from_millis(millis)) +} + +fn prune_asset_cache( + cache: &mut HashMap, + now: Instant, + ttl: Duration, + max_entries: usize, +) { + cache.retain(|_, entry| match entry { + AssetCacheEntry::Ready { fetched_at, .. } => now.duration_since(*fetched_at) < ttl, + AssetCacheEntry::Fetching(_) => true, + }); - pub fn spawn_sync(asset_db: Arc>) -> thread::JoinHandle<()> { - thread::spawn(move || loop { - if let Err(e) = asset_db.write().unwrap().fs_sync() { - error!("registry fs_sync failed: {:?}", e); + while cache.len() >= max_entries { + let oldest = cache + .iter() + .filter_map(|(asset_id, entry)| match entry { + AssetCacheEntry::Ready { fetched_at, .. } => Some((*asset_id, *fetched_at)), + AssetCacheEntry::Fetching(_) => None, + }) + .min_by_key(|(_, fetched_at)| *fetched_at) + .map(|(asset_id, _)| asset_id); + match oldest { + Some(asset_id) => { + cache.remove(&asset_id); } + None => break, + } + } +} - thread::sleep(Duration::from_secs(15)); - // TODO handle shutdowm +async fn decode_json_response( + mut response: reqwest::Response, + max_size: usize, +) -> std::result::Result { + if let Some(length) = response.content_length() { + if length > max_size as u64 { + return Err(RegistryError::InvalidResponse(format!( + "asset registry response exceeds {} bytes", + max_size + ))); + } + } + + let mut body = + Vec::with_capacity(response.content_length().unwrap_or(0).min(max_size as u64) as usize); + while let Some(chunk) = response.chunk().await.map_err(response_error)? { + let new_len = body.len().checked_add(chunk.len()).ok_or_else(|| { + RegistryError::InvalidResponse("asset registry response size overflow".to_string()) + })?; + if new_len > max_size { + return Err(RegistryError::InvalidResponse(format!( + "asset registry response exceeds {} bytes", + max_size + ))); + } + body.extend_from_slice(&chunk); + } + serde_json::from_slice(&body).map_err(|error| RegistryError::InvalidResponse(error.to_string())) +} + +fn request_error(error: reqwest::Error) -> RegistryError { + if error.is_timeout() { + RegistryError::Timeout(error.to_string()) + } else { + RegistryError::Transport(error.to_string()) + } +} + +fn response_error(error: reqwest::Error) -> RegistryError { + if error.is_timeout() { + RegistryError::Timeout(error.to_string()) + } else { + RegistryError::InvalidResponse(error.to_string()) + } +} + +#[derive(Deserialize)] +struct RegistryContractFields { + entity: JsonValue, + name: String, + precision: u8, + #[serde(default)] + ticker: Option, + version: u64, + #[serde(default)] + initial_issuer_pubkey: Option, + #[serde(default)] + issuer_pubkey: Option, + #[serde(flatten)] + extra: JsonMap, +} + +#[derive(Clone, Debug)] +pub struct RegistryContract { + pub entity: JsonValue, + pub name: String, + pub precision: u8, + pub ticker: Option, + pub version: u64, + pub initial_issuer_pubkey: Option, + pub issuer_pubkey: Option, + pub extra: JsonMap, + ticker_was_present: bool, + initial_issuer_pubkey_was_present: bool, + issuer_pubkey_was_present: bool, +} + +impl<'de> serde::Deserialize<'de> for RegistryContract { + fn deserialize(deserializer: D) -> std::result::Result + where + D: serde::Deserializer<'de>, + { + let value = ::deserialize(deserializer)?; + let object = value + .as_object() + .ok_or_else(|| serde::de::Error::custom("registry contract must be an object"))?; + let ticker_was_present = object.contains_key("ticker"); + let initial_issuer_pubkey_was_present = object.contains_key("initial_issuer_pubkey"); + let issuer_pubkey_was_present = object.contains_key("issuer_pubkey"); + let fields: RegistryContractFields = + serde_json::from_value(value).map_err(serde::de::Error::custom)?; + + Ok(Self { + entity: fields.entity, + name: fields.name, + precision: fields.precision, + ticker: fields.ticker, + version: fields.version, + initial_issuer_pubkey: fields.initial_issuer_pubkey, + issuer_pubkey: fields.issuer_pubkey, + extra: fields.extra, + ticker_was_present, + initial_issuer_pubkey_was_present, + issuer_pubkey_was_present, }) } } +impl serde::Serialize for RegistryContract { + fn serialize(&self, serializer: S) -> std::result::Result + where + S: serde::Serializer, + { + let mut contract = self.extra.clone(); + contract.insert("entity".to_string(), self.entity.clone()); + contract.insert("name".to_string(), JsonValue::String(self.name.clone())); + contract.insert("precision".to_string(), self.precision.into()); + contract.insert("version".to_string(), self.version.into()); + + insert_optional_contract_field( + &mut contract, + "ticker", + &self.ticker, + self.ticker_was_present, + ); + insert_optional_contract_field( + &mut contract, + "initial_issuer_pubkey", + &self.initial_issuer_pubkey, + self.initial_issuer_pubkey_was_present, + ); + insert_optional_contract_field( + &mut contract, + "issuer_pubkey", + &self.issuer_pubkey, + self.issuer_pubkey_was_present, + ); + + serde::Serialize::serialize(&JsonValue::Object(contract), serializer) + } +} + +fn insert_optional_contract_field( + contract: &mut JsonMap, + name: &str, + value: &Option, + was_present: bool, +) { + if was_present || value.is_some() { + contract.insert( + name.to_string(), + value + .as_ref() + .map(|value| JsonValue::String(value.clone())) + .unwrap_or(JsonValue::Null), + ); + } +} + +#[derive(Serialize, Deserialize, Clone, Debug)] +pub struct RegistryIcon { + pub href: String, + #[serde(flatten)] + pub extra: JsonMap, +} + +#[derive(Serialize, Deserialize, Clone, Debug)] +pub struct RegistryAsset { + pub asset_id: AssetId, + pub contract: RegistryContract, + pub initial_issuer_pubkey: String, + pub initial_issuer_pubkey_source: String, + pub current_issuer_pubkey: String, + #[serde(default)] + pub issuer_pubkey_history: Vec, + pub mutable: JsonValue, + #[serde(default)] + pub admin: Option, + #[serde(default)] + pub icon: Option, + pub status: String, + pub created_at: String, + pub updated_at: String, + #[serde(flatten)] + pub extra: JsonMap, +} + #[derive(Serialize, Deserialize, Clone, Debug)] pub struct AssetMeta { #[serde(skip_serializing_if = "JsonValue::is_null")] @@ -128,185 +648,847 @@ pub struct AssetMeta { pub name: String, #[serde(skip_serializing_if = "Option::is_none")] pub ticker: Option, + pub registry: RegistryAsset, } impl AssetMeta { - fn domain(&self) -> Option<&str> { - self.entity["domain"].as_str() + pub fn from_registry_asset( + registry: RegistryAsset, + ) -> std::result::Result { + let contract = serde_json::to_value(®istry.contract) + .map_err(|error| RegistryError::InvalidResponse(error.to_string()))?; + Ok(Self { + contract, + entity: registry.contract.entity.clone(), + precision: registry.contract.precision, + name: registry.contract.name.clone(), + ticker: registry.contract.ticker.clone(), + registry, + }) } } -pub struct AssetSorting(AssetSortField, AssetSortDir); - -pub enum AssetSortField { - Name, - Domain, - Ticker, +#[derive(Debug)] +pub struct RegistryAssetList { + pub total_count: usize, + pub items: Vec, } -pub enum AssetSortDir { - Descending, - Ascending, + +#[derive(Clone, Debug, Default, Eq, PartialEq)] +pub struct AssetSearchFilters { + asset_id: Option, + domain: Option, + ticker: Option, + name: Option, + asset_type: Option, + category_tags: Vec, + trading_venue: Option, + created_after: Option, + updated_after: Option, } -impl AssetSorting { - fn as_comparator(self) -> Box cmp::Ordering> { - let sort_fn: Box cmp::Ordering> = match self.0 { - AssetSortField::Name => { - // Order by name first, use asset id as a tie breaker. the other sorting fields - // don't require this because they're guaranteed to be unique. - Box::new(|a, b| lc_cmp(&a.1.name, &b.1.name).then_with(|| a.0.cmp(b.0))) - } - AssetSortField::Domain => Box::new(|a, b| a.1.domain().cmp(&b.1.domain())), - AssetSortField::Ticker => Box::new(|a, b| lc_cmp_opt(&a.1.ticker, &b.1.ticker)), +impl AssetSearchFilters { + pub fn new( + created_after: Option, + updated_after: Option, + ) -> Result { + let filters = Self { + created_after, + updated_after, + ..Self::default() + }; + filters.validate()?; + Ok(filters) + } + + pub fn from_query_pairs(query: &[(String, String)]) -> Result { + let get_last = |name: &str| { + query + .iter() + .rev() + .find(|(key, _)| key == name) + .map(|(_, value)| value.clone()) + }; + let filters = Self { + asset_id: get_last("asset_id"), + domain: get_last("domain"), + ticker: get_last("ticker"), + name: get_last("name"), + asset_type: get_last("asset_type"), + category_tags: query + .iter() + .filter(|(key, _)| key == "category_tag") + .map(|(_, value)| value.clone()) + .collect(), + trading_venue: get_last("trading_venue"), + created_after: get_last("created_after"), + updated_after: get_last("updated_after"), }; + filters.validate()?; + Ok(filters) + } - match self.1 { - AssetSortDir::Ascending => sort_fn, - AssetSortDir::Descending => Box::new(move |a, b| sort_fn(a, b).reverse()), + fn append_query<'a>(&'a self, query: &mut Vec<(&'a str, String)>) { + append_optional_query(query, "asset_id", &self.asset_id); + append_optional_query(query, "domain", &self.domain); + append_optional_query(query, "ticker", &self.ticker); + append_optional_query(query, "name", &self.name); + append_optional_query(query, "asset_type", &self.asset_type); + for category_tag in &self.category_tags { + query.push(("category_tag", category_tag.clone())); } + append_optional_query(query, "trading_venue", &self.trading_venue); + append_optional_query(query, "created_after", &self.created_after); + append_optional_query(query, "updated_after", &self.updated_after); } - pub fn from_query_params(query: &HashMap) -> Result { - let field = match query.get("sort_field").map(String::as_str) { - None => AssetSortField::Ticker, - Some("name") => AssetSortField::Name, - Some("domain") => AssetSortField::Domain, - Some("ticker") => AssetSortField::Ticker, - _ => bail!("invalid sort field"), - }; + fn validate(&self) -> Result<()> { + if let Some(value) = &self.asset_id { + ensure!( + (1..=64).contains(&value.len()) + && value.bytes().all(|byte| byte.is_ascii_hexdigit()), + "invalid asset_id: expected 1 to 64 hexadecimal characters" + ); + } + if let Some(value) = &self.domain { + ensure!( + is_valid_registry_domain(value), + "invalid domain: expected a valid domain name" + ); + } + validate_optional_registry_text("ticker", self.ticker.as_deref(), 24)?; + validate_optional_registry_text("name", self.name.as_deref(), 255)?; + validate_optional_registry_enum( + "asset_type", + self.asset_type.as_deref(), + &["AMP_asset", "stablecoin", "security_token", "other"], + )?; + for value in &self.category_tags { + validate_registry_enum( + "category_tag", + value, + &["stablecoin", "bond", "fixed-income", "tokenized"], + )?; + } + validate_optional_registry_enum( + "trading_venue", + self.trading_venue.as_deref(), + &["sideswap", "bitfinex"], + )?; + validate_optional_registry_timestamp("created_after", self.created_after.as_deref())?; + validate_optional_registry_timestamp("updated_after", self.updated_after.as_deref())?; + Ok(()) + } +} - let dir = match query.get("sort_dir").map(String::as_str) { - None => AssetSortDir::Ascending, - Some("asc") => AssetSortDir::Ascending, - Some("desc") => AssetSortDir::Descending, - _ => bail!("invalid sort direction"), - }; +fn append_optional_query<'a>( + query: &mut Vec<(&'a str, String)>, + name: &'a str, + value: &Option, +) { + if let Some(value) = value { + query.push((name, value.clone())); + } +} - Ok(Self(field, dir)) +fn validate_optional_registry_text(name: &str, value: Option<&str>, max_len: usize) -> Result<()> { + if let Some(value) = value { + ensure!( + value.chars().count() <= max_len && !value.contains('\0'), + "invalid {}: expected at most {} characters without NUL", + name, + max_len + ); } + Ok(()) } -fn lc_cmp(a: &str, b: &str) -> cmp::Ordering { - a.to_lowercase().cmp(&b.to_lowercase()) +fn validate_optional_registry_enum( + name: &str, + value: Option<&str>, + allowed: &[&str], +) -> Result<()> { + if let Some(value) = value { + validate_registry_enum(name, value, allowed)?; + } + Ok(()) } -fn lc_cmp_opt(a: &Option, b: &Option) -> cmp::Ordering { - a.as_ref() - .map(|a| a.to_lowercase()) - .cmp(&b.as_ref().map(|b| b.to_lowercase())) + +fn validate_registry_enum(name: &str, value: &str, allowed: &[&str]) -> Result<()> { + ensure!( + allowed + .iter() + .any(|candidate| candidate.eq_ignore_ascii_case(value)), + "invalid {}: expected one of {}", + name, + allowed.join(", ") + ); + Ok(()) +} + +fn is_valid_registry_domain(value: &str) -> bool { + if !(3..=255).contains(&value.len()) || !value.is_ascii() { + return false; + } + let value = value.strip_suffix('.').unwrap_or(value); + let labels: Vec<&str> = value.split('.').collect(); + if labels.len() < 2 { + return false; + } + labels.iter().all(|label| is_valid_registry_domain_label(label)) + && labels + .last() + .and_then(|label| label.bytes().next()) + .map(|byte| byte.is_ascii_alphabetic()) + .unwrap_or(false) +} + +fn is_valid_registry_domain_label(label: &str) -> bool { + if label.is_empty() || label.len() > 63 { + return false; + } + let first = label.bytes().next().unwrap(); + let last = label.bytes().last().unwrap(); + first.is_ascii_alphanumeric() + && last.is_ascii_alphanumeric() + && label + .bytes() + .all(|byte| byte.is_ascii_alphanumeric() || byte == b'-') +} + +fn validate_optional_registry_timestamp(name: &str, value: Option<&str>) -> Result<()> { + if let Some(value) = value { + OffsetDateTime::parse(value, &Rfc3339) + .map_err(|_| format!("invalid {}: expected an RFC 3339 date-time", name))?; + } + Ok(()) +} + +#[derive(Deserialize, Debug)] +struct RegistryListResponse { + items: Vec, + page: usize, + page_size: usize, + total_count: Option, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum AssetSorting { + AssetIdAsc, + AssetIdDesc, + NameAsc, + NameDesc, + DomainAsc, + DomainDesc, + TickerAsc, + TickerDesc, + CreatedAtAsc, + CreatedAtDesc, + UpdatedAtAsc, + UpdatedAtDesc, +} + +impl AssetSorting { + pub fn as_str(self) -> &'static str { + match self { + Self::AssetIdAsc => "asset_id_asc", + Self::AssetIdDesc => "asset_id_desc", + Self::NameAsc => "name_asc", + Self::NameDesc => "name_desc", + Self::DomainAsc => "domain_asc", + Self::DomainDesc => "domain_desc", + Self::TickerAsc => "ticker_asc", + Self::TickerDesc => "ticker_desc", + Self::CreatedAtAsc => "created_at_asc", + Self::CreatedAtDesc => "created_at_desc", + Self::UpdatedAtAsc => "updated_at_asc", + Self::UpdatedAtDesc => "updated_at_desc", + } + } + + pub fn from_query_params(query: &HashMap) -> Result { + if let Some(sort) = query.get("sort") { + ensure!( + !query.contains_key("sort_field") && !query.contains_key("sort_dir"), + "cannot combine sort with sort_field or sort_dir" + ); + return match sort.as_str() { + "asset_id_asc" => Ok(Self::AssetIdAsc), + "asset_id_desc" => Ok(Self::AssetIdDesc), + "name_asc" => Ok(Self::NameAsc), + "name_desc" => Ok(Self::NameDesc), + "domain_asc" => Ok(Self::DomainAsc), + "domain_desc" => Ok(Self::DomainDesc), + "ticker_asc" => Ok(Self::TickerAsc), + "ticker_desc" => Ok(Self::TickerDesc), + "created_at_asc" => Ok(Self::CreatedAtAsc), + "created_at_desc" => Ok(Self::CreatedAtDesc), + "updated_at_asc" => Ok(Self::UpdatedAtAsc), + "updated_at_desc" => Ok(Self::UpdatedAtDesc), + _ => bail!("invalid asset registry sort"), + }; + } + + let field = query + .get("sort_field") + .map(String::as_str) + .unwrap_or("ticker"); + let direction = query.get("sort_dir").map(String::as_str).unwrap_or("asc"); + match (field, direction) { + ("name", "asc") => Ok(Self::NameAsc), + ("name", "desc") => Ok(Self::NameDesc), + ("domain", "asc") => Ok(Self::DomainAsc), + ("domain", "desc") => Ok(Self::DomainDesc), + ("ticker", "asc") => Ok(Self::TickerAsc), + ("ticker", "desc") => Ok(Self::TickerDesc), + ("name" | "domain" | "ticker", _) => bail!("invalid sort direction"), + _ => bail!("invalid sort field"), + } + } } #[cfg(test)] mod tests { use super::*; - use tempfile::TempDir; + use std::io::{Read, Write}; + use std::net::TcpListener; + use std::str::FromStr; + use std::sync::mpsc; + use std::thread; const ASSET_ID_A: &str = "0000000000000000000000000000000000000000000000000000000000000001"; const ASSET_ID_B: &str = "0000000000000000000000000000000000000000000000000000000000000002"; - fn asset_id(id: &str) -> AssetId { - AssetId::from_str(id).unwrap() + fn search_filters(pairs: &[(&str, &str)]) -> AssetSearchFilters { + AssetSearchFilters::from_query_pairs( + &pairs + .iter() + .map(|(key, value)| (key.to_string(), value.to_string())) + .collect::>(), + ) + .unwrap() } - fn asset_file(dir: &TempDir, id: &str) -> path::PathBuf { - dir.path() - .join(&id[..DIR_PARTITION_LEN]) - .join(format!("{}.json", id)) + fn asset_response(asset_id: &str, name: &str, ticker: Option<&str>) -> JsonValue { + json!({ + "asset_id": asset_id, + "contract": { + "entity": {"domain": "example.com"}, + "name": name, + "precision": 8, + "ticker": ticker, + "version": 1, + "custom_contract_field": "preserved" + }, + "initial_issuer_pubkey": format!("02{}", "11".repeat(32)), + "initial_issuer_pubkey_source": "contract", + "current_issuer_pubkey": format!("02{}", "11".repeat(32)), + "issuer_pubkey_history": [], + "mutable": {"category_tags": ["stablecoin"], "custom": {"website": "https://example.com"}}, + "admin": {"featured": true}, + "icon": {"href": format!("/v2/assets/{}/icon/{}.png", asset_id, "22".repeat(32))}, + "status": "active", + "created_at": "2026-01-01T00:00:00Z", + "updated_at": "2026-01-02T00:00:00Z", + "future_field": {"preserved": true} + }) } - fn write_asset(dir: &TempDir, id: &str, name: &str, ticker: &str) { - let partition_dir = dir.path().join(&id[..DIR_PARTITION_LEN]); - fs::create_dir_all(&partition_dir).unwrap(); - fs::write( - partition_dir.join(format!("{}.json", id)), - format!( - r#"{{"contract":null,"entity":{{"domain":"example.com"}},"precision":8,"name":"{}","ticker":"{}"}}"#, - name, ticker - ), + fn mock_server( + responses: Vec<(u16, JsonValue)>, + ) -> (Url, mpsc::Receiver, thread::JoinHandle<()>) { + mock_server_with_delays( + responses + .into_iter() + .map(|(status, body)| (status, body, Duration::ZERO)) + .collect(), ) - .unwrap(); } - #[test] - fn fs_sync_loads_assets_from_disk() { - let dir = tempfile::tempdir().unwrap(); - write_asset(&dir, ASSET_ID_A, "Asset A", "AAA"); + fn mock_server_with_delays( + responses: Vec<(u16, JsonValue, Duration)>, + ) -> (Url, mpsc::Receiver, thread::JoinHandle<()>) { + let listener = TcpListener::bind("127.0.0.1:0").unwrap(); + let addr = listener.local_addr().unwrap(); + let (request_tx, request_rx) = mpsc::channel(); + let thread = thread::spawn(move || { + for (status, body, delay) in responses { + let (mut stream, _) = listener.accept().unwrap(); + stream + .set_read_timeout(Some(Duration::from_secs(1))) + .unwrap(); + let mut request = vec![0u8; 8192]; + let len = stream.read(&mut request).unwrap(); + request.truncate(len); + request_tx.send(String::from_utf8(request).unwrap()).ok(); + thread::sleep(delay); - let mut registry = AssetRegistry::new(dir.path().to_path_buf()); - registry.fs_sync().unwrap(); + let body = serde_json::to_string(&body).unwrap(); + let reason = match status { + 200 => "OK", + 404 => "Not Found", + 503 => "Service Unavailable", + _ => "Error", + }; + let response = format!( + "HTTP/1.1 {} {}\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", + status, + reason, + body.len(), + body + ); + stream.write_all(response.as_bytes()).unwrap(); + } + }); + ( + Url::parse(&format!("http://{}/api", addr)).unwrap(), + request_rx, + thread, + ) + } + + fn mock_unbounded_body(body: Vec) -> (Url, thread::JoinHandle<()>) { + let listener = TcpListener::bind("127.0.0.1:0").unwrap(); + let addr = listener.local_addr().unwrap(); + let thread = thread::spawn(move || { + let (mut stream, _) = listener.accept().unwrap(); + let mut request = vec![0u8; 8192]; + let _ = stream.read(&mut request); + let header = + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nConnection: close\r\n\r\n"; + let _ = stream.write_all(header.as_bytes()); + let _ = stream.write_all(&body); + }); + (Url::parse(&format!("http://{}/api", addr)).unwrap(), thread) + } + + fn mock_redirect() -> (Url, thread::JoinHandle<()>) { + let listener = TcpListener::bind("127.0.0.1:0").unwrap(); + let addr = listener.local_addr().unwrap(); + let thread = thread::spawn(move || { + let (mut stream, _) = listener.accept().unwrap(); + let mut request = vec![0u8; 8192]; + let _ = stream.read(&mut request); + let response = format!( + "HTTP/1.1 302 Found\r\nLocation: /api/v2/assets/{}\r\nContent-Length: 0\r\nConnection: close\r\n\r\n", + ASSET_ID_A + ); + stream.write_all(response.as_bytes()).unwrap(); + }); + (Url::parse(&format!("http://{}/api", addr)).unwrap(), thread) + } + + #[tokio::test] + async fn get_asset_projects_legacy_fields_and_preserves_v2_data() { + let body = asset_response(ASSET_ID_A, "Asset A", None); + let original_contract = body["contract"].clone(); + let (url, requests, server) = mock_server(vec![(200, body)]); + let expected_icon = url + .join(&format!( + "/v2/assets/{}/icon/{}.png", + ASSET_ID_A, + "22".repeat(32) + )) + .unwrap() + .to_string(); + let client = RegistryClient::new(url).unwrap(); + let id = AssetId::from_str(ASSET_ID_A).unwrap(); + + let asset = client.get_asset(&id).await.unwrap().unwrap(); + let metadata = AssetMeta::from_registry_asset(asset).unwrap(); - let metadata = registry.get(&asset_id(ASSET_ID_A)).unwrap(); assert_eq!(metadata.name, "Asset A"); - assert_eq!(metadata.ticker.as_deref(), Some("AAA")); + assert_eq!(metadata.ticker, None); + assert_eq!(metadata.contract, original_contract); + assert_eq!( + serde_json::to_value(&metadata.registry.contract).unwrap(), + original_contract + ); + assert_eq!(metadata.contract["custom_contract_field"], "preserved"); + assert_eq!(metadata.registry.mutable["category_tags"][0], "stablecoin"); + assert_eq!(metadata.registry.extra["future_field"]["preserved"], true); + assert_eq!(metadata.registry.icon.as_ref().unwrap().href, expected_icon); + assert!(requests + .recv() + .unwrap() + .starts_with(&format!("GET /api/v2/assets/{} HTTP/1.1", ASSET_ID_A))); + server.join().unwrap(); } - #[test] - fn fs_sync_removes_assets_deleted_from_disk() { - let dir = tempfile::tempdir().unwrap(); - write_asset(&dir, ASSET_ID_A, "Asset A", "AAA"); - write_asset(&dir, ASSET_ID_B, "Asset B", "BBB"); + #[tokio::test] + async fn get_asset_maps_not_found_and_rejects_mismatched_id() { + let mismatched = asset_response(ASSET_ID_A, "Asset A", Some("AAA")); + let (url, _, server) = mock_server(vec![(404, json!({})), (200, mismatched)]); + let client = RegistryClient::new(url).unwrap(); + let id_a = AssetId::from_str(ASSET_ID_A).unwrap(); + let id_b = AssetId::from_str(ASSET_ID_B).unwrap(); + + assert!(client.get_asset(&id_a).await.unwrap().is_none()); + assert!(matches!( + client.get_asset(&id_b).await, + Err(RegistryError::InvalidResponse(_)) + )); + server.join().unwrap(); + } - let mut registry = AssetRegistry::new(dir.path().to_path_buf()); - registry.fs_sync().unwrap(); - assert_eq!(registry.assets_cache.len(), 2); + #[tokio::test] + async fn asset_cache_hits_expires_and_coalesces_concurrent_misses() { + let body = asset_response(ASSET_ID_A, "Asset A", Some("AAA")); + let (url, requests, server) = + mock_server_with_delays(vec![(200, body, Duration::from_millis(40))]); + let client = RegistryClient::new(url).unwrap(); + let id = AssetId::from_str(ASSET_ID_A).unwrap(); + + let (first, second) = tokio::join!(client.get_asset(&id), client.get_asset(&id)); + assert!(first.unwrap().is_some()); + assert!(second.unwrap().is_some()); + assert!(client.get_asset(&id).await.unwrap().is_some()); + assert!(requests.recv().is_ok()); + server.join().unwrap(); + + let first = asset_response(ASSET_ID_A, "Asset A", Some("AAA")); + let second = asset_response(ASSET_ID_A, "Updated Asset A", Some("AAA")); + let (url, requests, server) = mock_server(vec![(200, first), (200, second)]); + let client = RegistryClient::with_options( + url, + Duration::from_secs(1), + Duration::from_secs(1), + Duration::from_millis(10), + 10, + 2, + ) + .unwrap(); + + assert_eq!( + client.get_asset(&id).await.unwrap().unwrap().contract.name, + "Asset A" + ); + tokio::time::sleep(Duration::from_millis(20)).await; + assert_eq!( + client.get_asset(&id).await.unwrap().unwrap().contract.name, + "Updated Asset A" + ); + assert!(requests.recv().is_ok()); + assert!(requests.recv().is_ok()); + server.join().unwrap(); + } + + #[tokio::test] + async fn client_rejects_excess_concurrent_requests() { + let body = asset_response(ASSET_ID_A, "Asset A", Some("AAA")); + let (url, requests, server) = + mock_server_with_delays(vec![(200, body, Duration::from_millis(100))]); + let client = RegistryClient::with_options( + url, + Duration::from_secs(1), + Duration::from_secs(1), + Duration::from_secs(1), + 10, + 1, + ) + .unwrap(); + let id_a = AssetId::from_str(ASSET_ID_A).unwrap(); + let id_b = AssetId::from_str(ASSET_ID_B).unwrap(); + let first_client = client.clone(); + let first = tokio::spawn(async move { first_client.get_asset(&id_a).await }); + + loop { + match requests.try_recv() { + Ok(_) => break, + Err(mpsc::TryRecvError::Empty) => { + tokio::time::sleep(Duration::from_millis(1)).await + } + Err(error) => panic!("mock registry stopped early: {}", error), + } + } + assert!(matches!( + client.get_asset(&id_b).await, + Err(RegistryError::Overloaded(_)) + )); + assert!(first.await.unwrap().unwrap().is_some()); + server.join().unwrap(); + } - fs::remove_file(asset_file(&dir, ASSET_ID_A)).unwrap(); - registry.fs_sync().unwrap(); + #[tokio::test] + async fn client_limits_streamed_response_bodies() { + let (url, server) = mock_unbounded_body(vec![b' '; REGISTRY_MAX_ASSET_RESPONSE_SIZE + 1]); + let client = RegistryClient::new(url).unwrap(); + let id = AssetId::from_str(ASSET_ID_A).unwrap(); - assert!(registry.get(&asset_id(ASSET_ID_A)).is_none()); - assert!(registry.get(&asset_id(ASSET_ID_B)).is_some()); - assert_eq!(registry.assets_cache.len(), 1); + assert!(matches!( + client.get_asset(&id).await, + Err(RegistryError::InvalidResponse(message)) + if message.contains("exceeds") + )); + server.join().unwrap(); + } + + #[tokio::test] + async fn client_rejects_redirects() { + let (url, server) = mock_redirect(); + let client = RegistryClient::new(url).unwrap(); + let id = AssetId::from_str(ASSET_ID_A).unwrap(); + + assert!(matches!( + client.get_asset(&id).await, + Err(RegistryError::HttpStatus(302)) + )); + server.join().unwrap(); } #[test] - fn fs_sync_reuses_cached_metadata_when_file_is_unchanged() { - let dir = tempfile::tempdir().unwrap(); - write_asset(&dir, ASSET_ID_A, "Asset A", "AAA"); + fn registry_timeout_environment_values_are_milliseconds() { + assert_eq!( + parse_registry_timeout(REGISTRY_CONNECT_TIMEOUT_ENV, "2500").unwrap(), + Duration::from_millis(2500) + ); + assert!(parse_registry_timeout(REGISTRY_CONNECT_TIMEOUT_ENV, "0").is_err()); + assert!(parse_registry_timeout(REGISTRY_REQUEST_TIMEOUT_ENV, "invalid").is_err()); + } - let mut registry = AssetRegistry::new(dir.path().to_path_buf()); - registry.fs_sync().unwrap(); + #[tokio::test] + async fn list_assets_translates_unaligned_offsets_across_pages() { + let first = json!({ + "items": [ + asset_response(ASSET_ID_A, "Asset A", Some("AAA")), + asset_response(ASSET_ID_B, "Asset B", Some("BBB")) + ], + "page": 2, + "page_size": 2, + "total_count": 5, + "total_pages": 3 + }); + let second = json!({ + "items": [asset_response(ASSET_ID_B, "Asset B", Some("BBB"))], + "page": 3, + "page_size": 2, + "total_count": 5, + "total_pages": 3 + }); + let (url, requests, server) = mock_server(vec![(200, first), (200, second)]); + let client = RegistryClient::new(url).unwrap(); - let id = asset_id(ASSET_ID_A); - let modified = fs::metadata(asset_file(&dir, ASSET_ID_A)) - .unwrap() - .modified() + let result = client + .list_assets( + 3, + 2, + AssetSorting::UpdatedAtAsc, + &search_filters(&[ + ("asset_id", "aB12"), + ("domain", "Example.com"), + ("ticker", "EXM"), + ("name", "Example"), + ("asset_type", "AMP_asset"), + ("category_tag", "stablecoin"), + ("category_tag", "bond"), + ("trading_venue", "sideswap"), + ("created_after", "2026-01-01T00:00:00Z"), + ("updated_after", "2026-02-01T12:30:00-05:00"), + ]), + ) + .await .unwrap(); - registry.assets_cache.insert( - id, - ( - modified, - AssetMeta { - contract: JsonValue::Null, - entity: json!({"domain": "cached.example.com"}), - precision: 0, - name: "Cached Asset".to_string(), - ticker: Some("CACHE".to_string()), - }, - ), - ); - registry.fs_sync().unwrap(); + assert_eq!(result.total_count, 5); + assert_eq!(result.items.len(), 2); + let first_request = requests.recv().unwrap(); + let second_request = requests.recv().unwrap(); + assert!(first_request.contains("page=2")); + assert!(first_request.contains("page_size=2")); + assert!(first_request.contains("sort=updated_at_asc")); + assert!(first_request.contains("asset_id=aB12")); + assert!(first_request.contains("domain=Example.com")); + assert!(first_request.contains("ticker=EXM")); + assert!(first_request.contains("name=Example")); + assert!(first_request.contains("asset_type=AMP_asset")); + assert!(first_request.contains("category_tag=stablecoin")); + assert!(first_request.contains("category_tag=bond")); + assert!(first_request.contains("trading_venue=sideswap")); + assert!(first_request.contains("created_after=2026-01-01T00%3A00%3A00Z")); + assert!(first_request.contains("updated_after=2026-02-01T12%3A30%3A00-05%3A00")); + assert!(second_request.contains("page=3")); + assert!(second_request.contains("sort=updated_at_asc")); + assert!(second_request.contains("name=Example")); + assert!(second_request.contains("category_tag=stablecoin")); + assert!(second_request.contains("category_tag=bond")); + assert!(second_request.contains("created_after=2026-01-01T00%3A00%3A00Z")); + assert!(second_request.contains("updated_after=2026-02-01T12%3A30%3A00-05%3A00")); + server.join().unwrap(); + } + + #[tokio::test] + async fn list_assets_requires_total_count() { + let body = json!({ + "items": [], + "page": 1, + "page_size": 25, + "total_count": null, + "total_pages": null + }); + let (url, _, server) = mock_server(vec![(200, body)]); + let client = RegistryClient::new(url).unwrap(); + + assert!(matches!( + client + .list_assets( + 0, + 25, + AssetSorting::TickerAsc, + &AssetSearchFilters::default() + ) + .await, + Err(RegistryError::InvalidResponse(_)) + )); + server.join().unwrap(); + } + + #[tokio::test] + async fn zero_limit_returns_only_the_total_count() { + let body = json!({ + "items": [asset_response(ASSET_ID_A, "Asset A", Some("AAA"))], + "page": 1, + "page_size": 1, + "total_count": 5, + "total_pages": 5 + }); + let (url, requests, server) = mock_server(vec![(200, body)]); + let client = RegistryClient::new(url).unwrap(); - let metadata = registry.get(&id).unwrap(); - assert_eq!(metadata.name, "Cached Asset"); - assert_eq!(metadata.ticker.as_deref(), Some("CACHE")); + let result = client + .list_assets( + usize::MAX, + 0, + AssetSorting::TickerAsc, + &AssetSearchFilters::default(), + ) + .await + .unwrap(); + assert_eq!(result.total_count, 5); + assert!(result.items.is_empty()); + let request = requests.recv().unwrap(); + assert!(request.contains("page=1")); + assert!(request.contains("page_size=1")); + server.join().unwrap(); + } + + #[tokio::test] + async fn pagination_rejects_overflow_without_an_http_request() { + let client = RegistryClient::new(Url::parse("http://127.0.0.1:1/").unwrap()).unwrap(); + assert!(matches!( + client + .list_assets( + usize::MAX, + 1, + AssetSorting::TickerAsc, + &AssetSearchFilters::default(), + ) + .await, + Err(RegistryError::InvalidRequest(_)) + )); + } + + #[tokio::test] + async fn client_distinguishes_http_status_and_timeout() { + let (url, _, server) = mock_server(vec![(503, json!({"detail": "unavailable"}))]); + let client = RegistryClient::new(url).unwrap(); + let id = AssetId::from_str(ASSET_ID_A).unwrap(); + assert!(matches!( + client.get_asset(&id).await, + Err(RegistryError::HttpStatus(503)) + )); + server.join().unwrap(); + + let listener = TcpListener::bind("127.0.0.1:0").unwrap(); + let addr = listener.local_addr().unwrap(); + let server = thread::spawn(move || { + let (_stream, _) = listener.accept().unwrap(); + thread::sleep(Duration::from_millis(100)); + }); + let client = RegistryClient::with_timeouts( + Url::parse(&format!("http://{}/", addr)).unwrap(), + Duration::from_millis(20), + Duration::from_millis(20), + ) + .unwrap(); + assert!(matches!( + client.get_asset(&id).await, + Err(RegistryError::Timeout(_)) + )); + server.join().unwrap(); } #[test] - fn fs_sync_refreshes_metadata_when_file_is_modified() { - let dir = tempfile::tempdir().unwrap(); - write_asset(&dir, ASSET_ID_A, "Asset A", "AAA"); + fn sorting_supports_legacy_and_native_parameters() { + let mut query = HashMap::new(); + query.insert("sort_field".to_string(), "domain".to_string()); + query.insert("sort_dir".to_string(), "desc".to_string()); + assert_eq!( + AssetSorting::from_query_params(&query).unwrap(), + AssetSorting::DomainDesc + ); - let mut registry = AssetRegistry::new(dir.path().to_path_buf()); - registry.fs_sync().unwrap(); + let mut query = HashMap::new(); + query.insert("sort".to_string(), "updated_at_desc".to_string()); + assert_eq!( + AssetSorting::from_query_params(&query).unwrap(), + AssetSorting::UpdatedAtDesc + ); + query.insert("sort_dir".to_string(), "asc".to_string()); + assert!(AssetSorting::from_query_params(&query).is_err()); - let id = asset_id(ASSET_ID_A); - registry.assets_cache.get_mut(&id).unwrap().0 = SystemTime::UNIX_EPOCH; - write_asset(&dir, ASSET_ID_A, "Updated Asset", "UPD"); + let mut query = HashMap::new(); + query.insert("sort".to_string(), "created_at_asc".to_string()); + assert_eq!( + AssetSorting::from_query_params(&query).unwrap(), + AssetSorting::CreatedAtAsc + ); + query.insert("sort".to_string(), "updated_at_asc".to_string()); + assert_eq!( + AssetSorting::from_query_params(&query).unwrap(), + AssetSorting::UpdatedAtAsc + ); + query.insert("sort".to_string(), "asset_id_desc".to_string()); + assert_eq!( + AssetSorting::from_query_params(&query).unwrap(), + AssetSorting::AssetIdDesc + ); + } - registry.fs_sync().unwrap(); + #[test] + fn search_filters_validate_the_openapi_constraints() { + let filters = search_filters(&[ + ("asset_id", "aB12"), + ("domain", "Sub.Example.com."), + ("ticker", "EXM"), + ("name", "Example Asset"), + ("asset_type", "amp_ASSET"), + ("category_tag", "StableCoin"), + ("category_tag", "fixed-income"), + ("trading_venue", "SideSwap"), + ("created_after", "2026-01-01T00:00:00Z"), + ("updated_after", "2026-02-01T12:30:00-05:00"), + ]); + assert_eq!(filters.category_tags.len(), 2); - let metadata = registry.get(&id).unwrap(); - assert_eq!(metadata.name, "Updated Asset"); - assert_eq!(metadata.ticker.as_deref(), Some("UPD")); + let invalid = [ + ("asset_id", "not-hex"), + ("domain", "invalid"), + ("ticker", "1234567890123456789012345"), + ("name", "contains\0nul"), + ("asset_type", "invalid"), + ("category_tag", "invalid"), + ("trading_venue", "invalid"), + ("created_after", "2026-01-01"), + ("updated_after", "2026-01-01T00:00:00"), + ]; + for (name, value) in invalid { + let query = vec![(name.to_string(), value.to_string())]; + assert!( + AssetSearchFilters::from_query_pairs(&query).is_err(), + "{} should reject {:?}", + name, + value + ); + } } } diff --git a/src/new_index/mod.rs b/src/new_index/mod.rs index 8d2a734ff..9d13e877e 100644 --- a/src/new_index/mod.rs +++ b/src/new_index/mod.rs @@ -12,6 +12,8 @@ pub use self::db::{DBRow, DB}; pub use self::fetch::{BlockEntry, FetchFrom}; pub use self::mempool::Mempool; pub use self::query::Query; +#[cfg(feature = "liquid")] +pub use self::query::{AssetLookup, AssetRegistryStatus}; pub use self::schema::{ compute_script_hash, parse_hash, ChainQuery, FundingInfo, GetAmountVal, Indexer, ScriptStats, SpendingInfo, SpendingInput, Store, TxHistoryInfo, TxHistoryKey, TxHistoryRow, Utxo, diff --git a/src/new_index/query.rs b/src/new_index/query.rs index 86dada56e..b202d7efc 100644 --- a/src/new_index/query.rs +++ b/src/new_index/query.rs @@ -16,7 +16,10 @@ use hyper::body::Bytes as BodyBytes; #[cfg(feature = "liquid")] use crate::{ chain::AssetId, - elements::{ebcompact::TxidCompat, lookup_asset, AssetRegistry, AssetSorting, LiquidAsset}, + elements::{ + ebcompact::TxidCompat, lookup_asset, AssetMeta, AssetSearchFilters, AssetSorting, + LiquidAsset, RegistryClient, RegistryError, + }, }; const FEE_ESTIMATES_TTL: u64 = 60; // seconds @@ -35,7 +38,22 @@ pub struct Query { cached_relayfee: RwLock>, cached_block_template: BlockTemplateCache, #[cfg(feature = "liquid")] - asset_db: Option>>, + asset_registry: Option>, +} + +#[cfg(feature = "liquid")] +#[derive(Debug)] +pub enum AssetRegistryStatus { + NotRequested, + Available, + NotFound, + Unavailable(RegistryError), +} + +#[cfg(feature = "liquid")] +pub struct AssetLookup { + pub asset: Option, + pub registry_status: AssetRegistryStatus, } impl Query { @@ -283,14 +301,14 @@ impl Query { mempool: Arc>, daemon: Arc, config: Arc, - asset_db: Option>>, + asset_registry: Option>, ) -> Self { Query { chain, mempool, daemon, config, - asset_db, + asset_registry, cached_estimates: RwLock::new((HashMap::new(), None)), cached_relayfee: RwLock::new(None), cached_block_template: BlockTemplateCache::new(), @@ -299,31 +317,93 @@ impl Query { #[cfg(feature = "liquid")] #[trace] - pub fn lookup_asset(&self, asset_id: &AssetId) -> Result> { - lookup_asset(&self, self.asset_db.as_ref(), asset_id, None) + pub fn lookup_asset_local(&self, asset_id: &AssetId) -> Result> { + lookup_asset(self, asset_id, None) + } + + #[cfg(feature = "liquid")] + #[trace] + pub async fn lookup_asset(&self, asset_id: &AssetId) -> Result { + let mut asset = match self.lookup_asset_local(asset_id)? { + Some(asset) => asset, + None => { + return Ok(AssetLookup { + asset: None, + registry_status: AssetRegistryStatus::NotRequested, + }) + } + }; + + if !matches!(asset, LiquidAsset::Issued(_)) { + return Ok(AssetLookup { + asset: Some(asset), + registry_status: AssetRegistryStatus::NotRequested, + }); + } + + let registry = match &self.asset_registry { + Some(registry) => registry, + None => { + return Ok(AssetLookup { + asset: Some(asset), + registry_status: AssetRegistryStatus::NotRequested, + }) + } + }; + + let registry_status = match registry.get_asset(asset_id).await { + Ok(Some(registry_asset)) => match AssetMeta::from_registry_asset(registry_asset) { + Ok(metadata) => { + if let LiquidAsset::Issued(issued) = &mut asset { + issued.meta = Some(metadata); + } + AssetRegistryStatus::Available + } + Err(error) => AssetRegistryStatus::Unavailable(error), + }, + Ok(None) => AssetRegistryStatus::NotFound, + Err(error) => AssetRegistryStatus::Unavailable(error), + }; + + Ok(AssetLookup { + asset: Some(asset), + registry_status, + }) } #[cfg(feature = "liquid")] #[trace] - pub fn list_registry_assets( + pub async fn list_registry_assets( &self, start_index: usize, limit: usize, sorting: AssetSorting, - ) -> Result<(usize, Vec)> { - let asset_db = match &self.asset_db { + filters: AssetSearchFilters, + ) -> std::result::Result<(usize, Vec), RegistryError> { + let registry = match &self.asset_registry { None => return Ok((0, vec![])), - Some(db) => db.read().unwrap(), + Some(registry) => registry, }; - let (total_num, results) = asset_db.list(start_index, limit, sorting); + + let page = registry + .list_assets(start_index, limit, sorting, &filters) + .await?; // Attach on-chain information alongside the registry metadata - let results = results - .into_iter() - .map(|(asset_id, metadata)| { - Ok(lookup_asset(&self, None, asset_id, Some(metadata))? - .chain_err(|| "missing registered asset")?) - }) - .collect::>>()?; - Ok((total_num, results)) + let total_count = page.total_count; + let mut results = Vec::with_capacity(page.items.len()); + for registry_asset in page.items { + let asset_id = registry_asset.asset_id; + let metadata = AssetMeta::from_registry_asset(registry_asset)?; + match lookup_asset(self, &asset_id, Some(metadata)) + .map_err(|error| RegistryError::LocalLookup(error.to_string()))? + { + Some(asset) => results.push(asset), + None => warn!( + "registered asset {} is not yet available in the local index", + asset_id + ), + } + } + Ok((total_count, results)) } } diff --git a/src/rest.rs b/src/rest.rs index 5493a0fc0..161a9c27c 100644 --- a/src/rest.rs +++ b/src/rest.rs @@ -6,6 +6,8 @@ use crate::config::Config; use crate::errors; use crate::new_index::{compute_script_hash, Query, SpendingInput, Utxo}; #[cfg(feature = "liquid")] +use crate::new_index::AssetRegistryStatus; +#[cfg(feature = "liquid")] use crate::util::optional_value_for_newer_blocks; use crate::util::{ create_socket, electrum_merkle, extract_tx_prevouts, get_innerscripts, get_tx_fee, has_prevout, @@ -33,7 +35,10 @@ use electrs_macros::trace; #[cfg(feature = "liquid")] use { - crate::elements::{ebcompact::*, peg::PegoutValue, AssetSorting, IssuanceValue}, + crate::elements::{ + ebcompact::*, peg::PegoutValue, AssetSearchFilters, AssetSorting, IssuanceValue, + RegistryError, + }, elements::{encode, secp256k1_zkp as zkp, AssetId}, }; @@ -566,6 +571,13 @@ fn spawn_conn( if let Some(ref origins) = config.cors { resp.headers_mut() .insert("Access-Control-Allow-Origin", origins.parse().unwrap()); + #[cfg(feature = "liquid")] + resp.headers_mut().insert( + "Access-Control-Expose-Headers", + "X-Asset-Registry-Status, X-Total-Results" + .parse() + .unwrap(), + ); } Ok::<_, hyper::Error>(resp) } @@ -678,7 +690,7 @@ impl Handle { } } -/// Whether `uri` addresses the block template endpoint, the one route handled on the async +/// Whether `uri` addresses the block template endpoint, one of the routes handled on the async /// runtime rather than on the blocking pool (see `handle_request`). Matched exactly the way /// the router below matches it, so the two cannot drift apart. fn is_block_template_request(method: &Method, uri: &hyper::Uri) -> bool { @@ -686,6 +698,19 @@ fn is_block_template_request(method: &Method, uri: &hyper::Uri) -> bool { *method == Method::GET && path.next() == Some("block-template") && path.next().is_none() } +#[cfg(feature = "liquid")] +fn is_asset_registry_request(method: &Method, uri: &hyper::Uri) -> bool { + if *method != Method::GET { + return false; + } + + let path: Vec<&str> = uri.path().split('/').skip(1).collect(); + matches!( + path.as_slice(), + ["assets", "registry"] | ["asset", _] | ["asset", _, "supply", "decimal"] + ) +} + /// Dispatch a request, keeping blocking work off the async worker threads. /// /// Almost every handler is synchronous: it reads RocksDB, and some of them (transaction @@ -695,9 +720,8 @@ fn is_block_template_request(method: &Method, uri: &hyper::Uri) -> bool { /// such as `GET /blocks/tip/height` stop being served. Moving them to the blocking pool /// keeps the runtime free to answer everything else. /// -/// The block template endpoint is the exception: it is genuinely asynchronous (concurrent -/// callers share one in-flight daemon fetch) and already does its own blocking work on the -/// blocking pool, so it stays on the runtime. +/// The block template and v2 asset registry endpoints are the exceptions: they are genuinely +/// asynchronous and already move or avoid blocking work, so they stay on the runtime. #[trace] async fn handle_request( method: Method, @@ -710,6 +734,11 @@ async fn handle_request( return handle_block_template_request(&query, &config).await; } + #[cfg(feature = "liquid")] + if is_asset_registry_request(&method, &uri) { + return handle_asset_registry_request(&uri, &query).await; + } + let path = uri.path().to_string(); tokio::task::spawn_blocking(move || handle_blocking_request(method, uri, body, &query, &config)) .await @@ -735,6 +764,102 @@ async fn handle_block_template_request( getblocktemplate_response(query.getblocktemplate().await) } +#[cfg(feature = "liquid")] +async fn handle_asset_registry_request( + uri: &hyper::Uri, + query: &Query, +) -> Result>, HttpError> { + let path: Vec<&str> = uri.path().split('/').skip(1).collect(); + let query_pairs = match uri.query() { + Some(value) => form_urlencoded::parse(value.as_bytes()) + .into_owned() + .collect::>(), + None => vec![], + }; + let query_params = query_pairs.iter().cloned().collect::>(); + + match path.as_slice() { + ["assets", "registry"] => { + let start_index: usize = query_params + .get("start_index") + .and_then(|n| n.parse().ok()) + .unwrap_or(0); + + let limit: usize = query_params + .get("limit") + .and_then(|n| n.parse().ok()) + .map(|n: usize| n.min(ASSETS_MAX_PER_PAGE)) + .unwrap_or(ASSETS_PER_PAGE); + + let sorting = AssetSorting::from_query_params(&query_params)?; + let filters = AssetSearchFilters::from_query_pairs(&query_pairs)?; + let (total_num, assets) = query + .list_registry_assets(start_index, limit, sorting, filters) + .await + .map_err(HttpError::from_registry_error)?; + + Ok(Response::builder() + // Disable caching because we don't currently support caching with query string params + .header("Cache-Control", "no-store") + .header("Content-Type", "application/json") + .header("X-Total-Results", total_num.to_string()) + .body(Full::new(Bytes::from(serde_json::to_string(&assets)?))) + .unwrap()) + } + ["asset", asset_str] => { + let asset_id = AssetId::from_str(asset_str)?; + let lookup = query.lookup_asset(&asset_id).await?; + let degraded = matches!( + &lookup.registry_status, + AssetRegistryStatus::Unavailable(_) + ); + let asset_entry = lookup + .asset + .ok_or_else(|| HttpError::not_found("Asset id not found".to_string()))?; + + let mut response = json_response_no_store(asset_entry, StatusCode::OK)?; + if degraded { + response.headers_mut().insert( + "X-Asset-Registry-Status", + "unavailable".parse().unwrap(), + ); + } + Ok(response) + } + ["asset", asset_str, "supply", "decimal"] => { + let asset_id = AssetId::from_str(asset_str)?; + let lookup = query.lookup_asset(&asset_id).await?; + let registry_error = match lookup.registry_status { + AssetRegistryStatus::Unavailable(error) => Some(error), + _ => None, + }; + let asset_entry = lookup + .asset + .ok_or_else(|| HttpError::not_found("Asset id not found".to_string()))?; + let supply = asset_entry + .supply() + .ok_or_else(|| HttpError::from("Asset supply is blinded".to_string()))?; + let precision = asset_entry.precision(); + + if precision > 0 { + http_message( + StatusCode::OK, + format_decimal_amount(supply, precision), + TTL_SHORT, + ) + } else if let Some(error) = registry_error { + Err(HttpError::from_registry_error(error)) + } else { + http_message(StatusCode::OK, supply.to_string(), TTL_SHORT) + } + } + _ => Err(HttpError::not_found(format!( + "endpoint does not exist {:?}", + uri.path() + ))), + } +} + /// The synchronous body of the router. Always invoked from the blocking pool by /// `handle_request`, never directly from an async worker thread. #[trace] @@ -1216,43 +1341,8 @@ fn handle_blocking_request( json_response(query.estimate_fee_map(), TTL_SHORT) } - // NOTE: `GET /block-template` is intercepted by `handle_request` before reaching - // here, because it is the only asynchronous handler. See `is_block_template_request`. - #[cfg(feature = "liquid")] - (&Method::GET, Some(&"assets"), Some(&"registry"), None, None, None) => { - let start_index: usize = query_params - .get("start_index") - .and_then(|n| n.parse().ok()) - .unwrap_or(0); - - let limit: usize = query_params - .get("limit") - .and_then(|n| n.parse().ok()) - .map(|n: usize| n.min(ASSETS_MAX_PER_PAGE)) - .unwrap_or(ASSETS_PER_PAGE); - - let sorting = AssetSorting::from_query_params(&query_params)?; - - let (total_num, assets) = query.list_registry_assets(start_index, limit, sorting)?; - - Ok(Response::builder() - // Disable caching because we don't currently support caching with query string params - .header("Cache-Control", "no-store") - .header("Content-Type", "application/json") - .header("X-Total-Results", total_num.to_string()) - .body(Full::new(Bytes::from(serde_json::to_string(&assets)?))) - .unwrap()) - } - - #[cfg(feature = "liquid")] - (&Method::GET, Some(&"asset"), Some(asset_str), None, None, None) => { - let asset_id = AssetId::from_str(asset_str)?; - let asset_entry = query - .lookup_asset(&asset_id)? - .ok_or_else(|| HttpError::not_found("Asset id not found".to_string()))?; - - json_response(asset_entry, TTL_SHORT) - } + // NOTE: asynchronous endpoints are intercepted by `handle_request` before reaching + // this synchronous router. See the route classifiers above. #[cfg(feature = "liquid")] (&Method::GET, Some(&"asset"), Some(asset_str), Some(&"txs"), None, None) => { @@ -1316,23 +1406,15 @@ fn handle_blocking_request( } #[cfg(feature = "liquid")] - (&Method::GET, Some(&"asset"), Some(asset_str), Some(&"supply"), param, None) => { + (&Method::GET, Some(&"asset"), Some(asset_str), Some(&"supply"), None, None) => { let asset_id = AssetId::from_str(asset_str)?; let asset_entry = query - .lookup_asset(&asset_id)? + .lookup_asset_local(&asset_id)? .ok_or_else(|| HttpError::not_found("Asset id not found".to_string()))?; - let supply = asset_entry .supply() .ok_or_else(|| HttpError::from("Asset supply is blinded".to_string()))?; - let precision = asset_entry.precision(); - - if param == Some(&"decimal") && precision > 0 { - let supply_dec = supply as f64 / 10u32.pow(precision.into()) as f64; - http_message(StatusCode::OK, supply_dec.to_string(), TTL_SHORT) - } else { - http_message(StatusCode::OK, supply.to_string(), TTL_SHORT) - } + http_message(StatusCode::OK, supply.to_string(), TTL_SHORT) } _ => Err(HttpError::not_found(format!( @@ -1358,6 +1440,31 @@ where .unwrap()) } +#[cfg(feature = "liquid")] +fn format_decimal_amount(amount: u64, precision: u8) -> String { + if precision == 0 { + return amount.to_string(); + } + + let precision = usize::from(precision); + let digits = amount.to_string(); + let (whole, fractional) = if digits.len() > precision { + let split = digits.len() - precision; + (digits[..split].to_string(), digits[split..].to_string()) + } else { + ( + "0".to_string(), + format!("{}{}", "0".repeat(precision - digits.len()), digits), + ) + }; + let fractional = fractional.trim_end_matches('0'); + if fractional.is_empty() { + whole + } else { + format!("{}.{}", whole, fractional) + } +} + fn json_response(value: T, ttl: u32) -> Result>, HttpError> { json_response_with_status(value, StatusCode::OK, ttl) } @@ -1527,6 +1634,23 @@ impl HttpError { fn forbidden(msg: String) -> Self { HttpError(StatusCode::FORBIDDEN, msg) } + + #[cfg(feature = "liquid")] + fn from_registry_error(error: RegistryError) -> Self { + let status = match &error { + RegistryError::InvalidRequest(_) => StatusCode::BAD_REQUEST, + RegistryError::HttpStatus(400 | 422) => StatusCode::BAD_REQUEST, + RegistryError::Timeout(_) => StatusCode::GATEWAY_TIMEOUT, + RegistryError::HttpStatus(429 | 503) + | RegistryError::Overloaded(_) => StatusCode::SERVICE_UNAVAILABLE, + RegistryError::LocalLookup(_) => StatusCode::INTERNAL_SERVER_ERROR, + RegistryError::InvalidBaseUrl(_) + | RegistryError::Transport(_) + | RegistryError::HttpStatus(_) + | RegistryError::InvalidResponse(_) => StatusCode::BAD_GATEWAY, + }; + HttpError(status, error.to_string()) + } } impl From for HttpError { @@ -1621,6 +1745,10 @@ impl From for HttpError { #[cfg(test)] mod tests { + #[cfg(feature = "liquid")] + use crate::elements::RegistryError; + #[cfg(feature = "liquid")] + use crate::rest::is_asset_registry_request; use crate::rest::{is_block_template_request, HttpError}; use crate::{errors, errors::ErrorKind}; use http_body_util::BodyExt; @@ -1628,20 +1756,83 @@ mod tests { use serde_json::Value; use std::collections::HashMap; + #[cfg(feature = "liquid")] + #[test] + fn registry_errors_map_to_gateway_statuses() { + assert_eq!( + HttpError::from_registry_error(RegistryError::Timeout("timeout".to_string())).0, + StatusCode::GATEWAY_TIMEOUT + ); + assert_eq!( + HttpError::from_registry_error(RegistryError::HttpStatus(503)).0, + StatusCode::SERVICE_UNAVAILABLE + ); + assert_eq!( + HttpError::from_registry_error(RegistryError::HttpStatus(400)).0, + StatusCode::BAD_REQUEST + ); + assert_eq!( + HttpError::from_registry_error(RegistryError::HttpStatus(422)).0, + StatusCode::BAD_REQUEST + ); + assert_eq!( + HttpError::from_registry_error(RegistryError::Overloaded("busy".to_string())).0, + StatusCode::SERVICE_UNAVAILABLE + ); + assert_eq!( + HttpError::from_registry_error(RegistryError::InvalidResponse("bad json".to_string())) + .0, + StatusCode::BAD_GATEWAY + ); + } + + #[cfg(feature = "liquid")] + #[test] + fn decimal_asset_amounts_are_formatted_without_overflow_or_rounding() { + assert_eq!(super::format_decimal_amount(1_500_000_000, 10), "0.15"); + assert_eq!( + super::format_decimal_amount(1, 18), + "0.000000000000000001" + ); + assert_eq!( + super::format_decimal_amount(u64::MAX, 18), + "18.446744073709551615" + ); + assert_eq!(super::format_decimal_amount(0, 18), "0"); + } + #[test] - fn block_template_is_the_only_route_kept_on_the_async_runtime() { + fn async_routes_are_kept_on_the_async_runtime() { let is_async = |method: Method, uri: &str| { - is_block_template_request(&method, &uri.parse::().unwrap()) + let uri = uri.parse::().unwrap(); + let is_async = is_block_template_request(&method, &uri); + #[cfg(feature = "liquid")] + let is_async = is_async || is_asset_registry_request(&method, &uri); + is_async }; assert!(is_async(Method::GET, "/block-template")); assert!(is_async(Method::GET, "/block-template?ignored=1")); - // Everything else must fall through to the blocking pool, including near-misses - // that the router itself would not match as the block template route. assert!(!is_async(Method::GET, "/block-template/")); assert!(!is_async(Method::GET, "/block-template/extra")); assert!(!is_async(Method::POST, "/block-template")); + + #[cfg(feature = "liquid")] + { + assert!(is_async(Method::GET, "/assets/registry")); + assert!(is_async(Method::GET, "/assets/registry?limit=5")); + assert!(is_async(Method::GET, "/asset/asset-id")); + assert!(is_async( + Method::GET, + "/asset/asset-id/supply/decimal" + )); + assert!(!is_async(Method::POST, "/assets/registry")); + assert!(!is_async(Method::GET, "/assets/registry/")); + assert!(!is_async(Method::GET, "/asset/asset-id/supply")); + assert!(!is_async(Method::GET, "/asset/asset-id/txs")); + } + assert!(!is_async(Method::GET, "/blocks/tip/height")); assert!(!is_async(Method::POST, "/tx")); } diff --git a/tests/common.rs b/tests/common.rs index f43d88b90..43b886f61 100644 --- a/tests/common.rs +++ b/tests/common.rs @@ -28,6 +28,10 @@ use electrs::{ rest, signal::Waiter, }; +#[cfg(feature = "liquid")] +use electrs::elements::RegistryClient; +#[cfg(feature = "liquid")] +use url::Url; pub struct TestRunner { config: Arc, @@ -44,6 +48,25 @@ pub struct TestRunner { impl TestRunner { pub fn new() -> Result { + Self::new_inner( + #[cfg(feature = "liquid")] + None, + None, + ) + } + + #[cfg(feature = "liquid")] + pub fn new_with_asset_registry( + asset_registry_url: Url, + cors: Option, + ) -> Result { + Self::new_inner(Some(asset_registry_url), cors) + } + + fn new_inner( + #[cfg(feature = "liquid")] asset_registry_url: Option, + cors: Option, + ) -> Result { let log = init_log(); // Setup the bitcoind/elementsd config @@ -109,7 +132,7 @@ impl TestRunner { address_search: true, index_unspendables: false, enable_mining_rest: true, - cors: None, + cors, precache_scripts: None, utxos_limit: 100, electrum_txs_limit: 100, @@ -119,7 +142,7 @@ impl TestRunner { zmq_addr: None, #[cfg(feature = "liquid")] - asset_db_path: None, // XXX + asset_registry_url, #[cfg(feature = "liquid")] parent_network: bitcoin::Network::Regtest, db_block_cache_mb: 8, @@ -183,13 +206,22 @@ impl TestRunner { ))); assert!(Mempool::update(&mempool, &daemon, &tip)?); + #[cfg(feature = "liquid")] + let asset_registry = config + .asset_registry_url + .clone() + .map(RegistryClient::new) + .transpose() + .chain_err(|| "failed creating test asset registry client")? + .map(Arc::new); + let query = Arc::new(Query::new( Arc::clone(&chain), Arc::clone(&mempool), Arc::clone(&daemon), Arc::clone(&config), #[cfg(feature = "liquid")] - None, // TODO + asset_registry, )); let salt_rwlock = Arc::new(RwLock::new(String::from("foobar"))); @@ -330,6 +362,18 @@ pub fn init_rest_tester() -> Result<(rest::Handle, net::SocketAddr, TestRunner)> wait_for_tcp(addr, "REST"); Ok((rest_server, addr, tester)) } + +#[cfg(feature = "liquid")] +pub fn init_rest_tester_with_asset_registry( + asset_registry_url: Url, + cors: Option, +) -> Result<(rest::Handle, net::SocketAddr, TestRunner)> { + let tester = TestRunner::new_with_asset_registry(asset_registry_url, cors)?; + let addr = tester.config.http_addr; + let rest_server = rest::start(Arc::clone(&tester.config), Arc::clone(&tester.query)); + wait_for_tcp(addr, "REST"); + Ok((rest_server, addr, tester)) +} pub fn init_electrum_tester() -> Result<(ElectrumRPC, net::SocketAddr, TestRunner)> { let tester = TestRunner::new()?; let addr = tester.config.electrum_rpc_addr; diff --git a/tests/rest.rs b/tests/rest.rs index 91100cf9c..765b289e8 100644 --- a/tests/rest.rs +++ b/tests/rest.rs @@ -1,9 +1,20 @@ use bitcoin::hashes::{sha256, Hash}; use bitcoin::hex::FromHex; -use serde_json::Value; -use std::collections::HashSet; +use serde_json::{json, Value}; +use std::collections::{HashMap, HashSet}; use std::net; +#[cfg(feature = "liquid")] +use std::io::{Read, Write}; +#[cfg(feature = "liquid")] +use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; +#[cfg(feature = "liquid")] +use std::sync::{Arc, Mutex}; +#[cfg(feature = "liquid")] +use std::thread; +#[cfg(feature = "liquid")] +use url::Url; + #[cfg(feature = "liquid")] use elementsd::bitcoincore_rpc::RpcApi; #[cfg(not(feature = "liquid"))] @@ -30,6 +41,150 @@ fn get_plain(rest_addr: net::SocketAddr, path: &str) -> Result { Ok(get(rest_addr, path)?.into_body().read_to_string()?) } +#[cfg(feature = "liquid")] +fn registry_asset_response(asset_id: &str) -> Value { + json!({ + "asset_id": asset_id, + "contract": { + "entity": {"domain": "example.com"}, + "name": "Registry Asset", + "precision": 8, + "ticker": "REG", + "version": 1, + "custom_contract_field": "preserved" + }, + "initial_issuer_pubkey": format!("02{}", "11".repeat(32)), + "initial_issuer_pubkey_source": "contract", + "current_issuer_pubkey": format!("02{}", "11".repeat(32)), + "issuer_pubkey_history": [], + "mutable": {"category_tags": ["stablecoin"]}, + "admin": {"featured": true}, + "icon": {"href": format!("/v2/assets/{}/icon/{}.png", asset_id, "22".repeat(32))}, + "status": "active", + "created_at": "2026-01-01T00:00:00Z", + "updated_at": "2026-01-02T00:00:00Z" + }) +} + +#[cfg(feature = "liquid")] +fn missing_registry_asset_response() -> Value { + registry_asset_response( + "1111111111111111111111111111111111111111111111111111111111111111", + ) +} + +#[cfg(feature = "liquid")] +fn start_asset_registry_mock( + expected_requests: usize, +) -> ( + Url, + Arc>>, + Arc, + Arc, + thread::JoinHandle<()>, +) { + let listener = net::TcpListener::bind("127.0.0.1:0").unwrap(); + let addr = listener.local_addr().unwrap(); + let asset_id = Arc::new(Mutex::new(None::)); + let available = Arc::new(AtomicBool::new(true)); + let request_count = Arc::new(AtomicUsize::new(0)); + let server_asset_id = Arc::clone(&asset_id); + let server_available = Arc::clone(&available); + let server_request_count = Arc::clone(&request_count); + let thread = thread::spawn(move || { + for _ in 0..expected_requests { + let (mut stream, _) = listener.accept().unwrap(); + let mut request = vec![0u8; 8192]; + let len = stream.read(&mut request).unwrap(); + let request = String::from_utf8(request[..len].to_vec()).unwrap(); + let path = request + .lines() + .next() + .and_then(|line| line.split_whitespace().nth(1)) + .unwrap(); + server_request_count.fetch_add(1, Ordering::SeqCst); + + let (status, reason, body) = if server_available.load(Ordering::SeqCst) { + let asset_id = server_asset_id.lock().unwrap().clone().unwrap(); + let asset = registry_asset_response(&asset_id); + let body = if path.starts_with("/api/v2/assets?") { + let url = Url::parse(&format!("http://registry.invalid{}", path)).unwrap(); + let query_pairs: Vec<(String, String)> = + url.query_pairs().into_owned().collect(); + let query: HashMap = + query_pairs.iter().cloned().collect(); + assert_eq!(query.get("asset_id").map(String::as_str), Some("aB12")); + assert_eq!( + query.get("domain").map(String::as_str), + Some("Example.com") + ); + assert_eq!(query.get("ticker").map(String::as_str), Some("EXM")); + assert_eq!( + query.get("name").map(String::as_str), + Some("Registry") + ); + assert_eq!( + query.get("asset_type").map(String::as_str), + Some("AMP_asset") + ); + let category_tags: Vec<&str> = query_pairs + .iter() + .filter(|(key, _)| key == "category_tag") + .map(|(_, value)| value.as_str()) + .collect(); + assert_eq!(category_tags, ["stablecoin", "bond"]); + assert_eq!( + query.get("trading_venue").map(String::as_str), + Some("sideswap") + ); + assert_eq!( + query.get("created_after").map(String::as_str), + Some("2026-01-01T00:00:00Z") + ); + assert_eq!( + query.get("updated_after").map(String::as_str), + Some("2026-02-01T12:30:00-05:00") + ); + assert_eq!( + query.get("sort").map(String::as_str), + Some("created_at_asc") + ); + json!({ + "items": [asset, missing_registry_asset_response()], + "page": 1, + "page_size": 25, + "total_count": 2, + "total_pages": 1 + }) + } else { + assert_eq!(path, format!("/api/v2/assets/{}", asset_id)); + asset + }; + (200, "OK", body) + } else { + (503, "Service Unavailable", json!({"detail": "unavailable"})) + }; + let body = serde_json::to_string(&body).unwrap(); + let response = format!( + "HTTP/1.1 {} {}\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", + status, + reason, + body.len(), + body + ); + stream.write_all(response.as_bytes()).unwrap(); + } + }); + + ( + Url::parse(&format!("http://{}/api", addr)).unwrap(), + asset_id, + available, + request_count, + thread, + ) +} + #[test] fn test_rest_tx() -> Result<()> { let (rest_handle, rest_addr, mut tester) = common::init_rest_tester().unwrap(); @@ -1478,6 +1633,175 @@ fn test_rest_liquid_unblinded_issuance() -> Result<()> { Ok(()) } +#[cfg(feature = "liquid")] +#[test] +fn test_rest_liquid_v2_asset_registry() -> Result<()> { + let (registry_url, registry_asset_id, registry_available, request_count, registry_thread) = + start_asset_registry_mock(4); + let registry_url_for_icon = registry_url.clone(); + let (rest_handle, rest_addr, mut tester) = + common::init_rest_tester_with_asset_registry(registry_url, Some("*".to_string()))?; + + let issuance = tester + .node_client() + .call::("issueasset", &[1.5.into(), 0.into(), false.into()])?; + tester.mine()?; + let asset_id = issuance["asset"].as_str().unwrap().to_string(); + *registry_asset_id.lock().unwrap() = Some(asset_id.clone()); + let expected_icon = registry_url_for_icon + .join(&format!( + "/v2/assets/{}/icon/{}.png", + asset_id, + "22".repeat(32) + )) + .unwrap() + .to_string(); + + let response = get(rest_addr, &format!("/asset/{}", asset_id))?; + assert_eq!( + response + .headers() + .get("cache-control") + .and_then(|value| value.to_str().ok()), + Some("no-store") + ); + let asset: Value = response.into_body().read_json()?; + assert_eq!(asset["name"], "Registry Asset"); + assert_eq!(asset["ticker"], "REG"); + assert_eq!(asset["precision"], 8); + assert_eq!(asset["contract"]["custom_contract_field"], "preserved"); + assert_eq!(asset["registry"]["status"], "active"); + assert_eq!(asset["registry"]["mutable"]["category_tags"][0], "stablecoin"); + assert_eq!(asset["registry"]["icon"]["href"], expected_icon); + + assert_eq!( + get_plain( + rest_addr, + &format!("/asset/{}/supply/decimal", asset_id) + )?, + "1.5" + ); + assert_eq!( + get_plain(rest_addr, &format!("/asset/{}/supply", asset_id))?, + "150000000" + ); + + let response = get( + rest_addr, + concat!( + "/assets/registry?asset_id=aB12&domain=Example.com&ticker=EXM&name=Registry", + "&asset_type=AMP_asset&category_tag=stablecoin&category_tag=bond", + "&trading_venue=sideswap&created_after=2026-01-01T00%3A00%3A00Z", + "&updated_after=2026-02-01T12%3A30%3A00-05%3A00&sort=created_at_asc" + ), + )?; + assert_eq!( + response + .headers() + .get("cache-control") + .and_then(|value| value.to_str().ok()), + Some("no-store") + ); + assert_eq!( + response + .headers() + .get("x-total-results") + .and_then(|value| value.to_str().ok()), + Some("2") + ); + let assets: Value = response.into_body().read_json()?; + assert_eq!(assets.as_array().unwrap().len(), 1); + assert_eq!(assets[0]["asset_id"], asset_id); + assert_eq!(assets[0]["registry"]["status"], "active"); + assert_eq!(assets[0]["registry"]["icon"]["href"], expected_icon); + + let response = ureq::get(&format!( + "http://{}/assets/registry?created_after=2026-01-01", + rest_addr + )) + .config() + .http_status_as_error(false) + .build() + .call()?; + assert_eq!(response.status(), 400); + assert_eq!( + response.into_body().read_to_string()?, + "invalid created_after: expected an RFC 3339 date-time" + ); + + let response = ureq::get(&format!( + "http://{}/assets/registry?updated_after=not-a-time", + rest_addr + )) + .config() + .http_status_as_error(false) + .build() + .call()?; + assert_eq!(response.status(), 400); + assert_eq!( + response.into_body().read_to_string()?, + "invalid updated_after: expected an RFC 3339 date-time" + ); + + let response = ureq::get(&format!( + "http://{}/assets/registry?asset_id=not-hex", + rest_addr + )) + .config() + .http_status_as_error(false) + .build() + .call()?; + assert_eq!(response.status(), 400); + assert_eq!( + response.into_body().read_to_string()?, + "invalid asset_id: expected 1 to 64 hexadecimal characters" + ); + + thread::sleep(std::time::Duration::from_millis(1100)); + registry_available.store(false, Ordering::SeqCst); + let response = get(rest_addr, &format!("/asset/{}", asset_id))?; + assert_eq!(response.status(), 200); + assert_eq!( + response + .headers() + .get("x-asset-registry-status") + .and_then(|value| value.to_str().ok()), + Some("unavailable") + ); + assert_eq!( + response + .headers() + .get("access-control-allow-origin") + .and_then(|value| value.to_str().ok()), + Some("*") + ); + assert_eq!( + response + .headers() + .get("access-control-expose-headers") + .and_then(|value| value.to_str().ok()), + Some("X-Asset-Registry-Status, X-Total-Results") + ); + let degraded_asset: Value = response.into_body().read_json()?; + assert!(degraded_asset.get("registry").is_none()); + assert!(degraded_asset.get("name").is_none()); + + let response = ureq::get(&format!( + "http://{}/asset/{}/supply/decimal", + rest_addr, asset_id + )) + .config() + .http_status_as_error(false) + .build() + .call()?; + assert_eq!(response.status(), 503); + + assert_eq!(request_count.load(Ordering::SeqCst), 4); + rest_handle.stop(); + registry_thread.join().unwrap(); + Ok(()) +} + #[cfg(feature = "liquid")] #[test] fn test_rest_liquid_asset_transfer() -> Result<()> {