diff --git a/internal/config_specification.toml b/internal/config_specification.toml index 9cabe37..7405b0e 100644 --- a/internal/config_specification.toml +++ b/internal/config_specification.toml @@ -63,3 +63,9 @@ name = "debug" type = "bool" doc = "Enable debug RPC endpoints (e.g. txchain inspection). Do not use in production." default = "false" + +[[param]] +name = "riften_ipfs_gateway" +type = "String" +doc = "Optional IPFS gateway prefix for the local riften-ipfs pinning node (e.g. 'http://127.0.0.1:3002/ipfs/'). When set, ipfs:// BCMR content is fetched from here first, before the public gateways — so content pinned on the riften node is loadable without waiting for public DHT propagation. Empty disables it." +default = "\"\".to_string()" diff --git a/src/bcmr/bcmrdownloader.rs b/src/bcmr/bcmrdownloader.rs index 0c2a110..c6f7034 100644 --- a/src/bcmr/bcmrdownloader.rs +++ b/src/bcmr/bcmrdownloader.rs @@ -31,6 +31,9 @@ use super::utilurl::get_url; // We are generous on timeout to allow for slow ipfs gateway const DOWNLOAD_TIMEOUT: Duration = Duration::from_secs(60); +// The local riften gateway is on the same host; if it hangs (vs. refusing) we want to fall +// back to the public gateways quickly rather than wait out the full DOWNLOAD_TIMEOUT. +const LOCAL_GATEWAY_TIMEOUT: Duration = Duration::from_secs(10); const IPFS_GATEWAYS: [&str; 2] = ["https://ipfs.io/ipfs/", "https://w3s.link/ipfs/"]; const MAX_BCMR_SIZE: usize = 1024 * 1024 * 100; // 100 MB @@ -38,21 +41,53 @@ const MAX_PARALLEL_DOWNLOADS: usize = 5; const WELL_KNOWN_PATH: &str = "/.well-known/bitcoin-cash-metadata-registry.json"; +/// Validate the post-`ipfs://` portion (a CID, optionally followed by a `/`-separated path) +/// before it is concatenated onto a gateway URL. The gateway prefix fixes the host, so the +/// host can't be swapped here, but `stripped` is untrusted on-chain data — reject path +/// traversal and URL-structure tricks so a crafted value can't escape the gateway's `/ipfs/` +/// path, inject request headers (CRLF), or smuggle a port/userinfo. +fn is_safe_ipfs_path(stripped: &str) -> bool { + let first = stripped.split('/').next().unwrap_or(""); + if first.is_empty() || first.starts_with('.') { + return false; + } + !(stripped.contains("..") + || stripped.contains("//") + || stripped.contains('@') + || stripped.contains(':') + || stripped.contains('\\') + || stripped.starts_with('/') + || stripped + .chars() + .any(|c| c.is_whitespace() || c.is_control())) +} + /// Resolve an on-chain BCMR URI into one or more candidate URLs to try in order. /// /// Handles three cases: -/// - `ipfs://` → a single random IPFS gateway URL. +/// - `ipfs://` → the local riften-ipfs gateway first (if configured), then a random +/// public IPFS gateway. Trying the riften node first means content pinned there is +/// loadable immediately, without waiting for public DHT propagation. The CID/path is +/// validated first (see `is_safe_ipfs_path`); an unsafe value yields no candidates. /// - Bare domain (authority only, no path) → BCMR v2 well-known URI first, then /// the bare root URL as a fallback for bug-compatibility with tokens that /// served the registry directly at `/`. /// - Anything with a path → used as-is (with `https://` prepended if missing). -fn resolve_uri_candidates(uri: &str) -> Vec { +fn resolve_uri_candidates(uri: &str, riften_gateway: &str) -> Vec { if let Some(stripped) = uri.strip_prefix("ipfs://") { - return vec![format!( + if !is_safe_ipfs_path(stripped) { + return Vec::new(); + } + let mut urls = Vec::new(); + if !riften_gateway.is_empty() { + urls.push(format!("{riften_gateway}{stripped}")); + } + urls.push(format!( "{}{}", IPFS_GATEWAYS.choose(&mut thread_rng()).unwrap(), stripped - )]; + )); + return urls; } let (scheme, rest) = match uri.split_once("://") { @@ -79,6 +114,9 @@ pub struct BCMRDownloader { db: SqlitePool, keep_running: Arc, download_task: Option>, + /// IPFS gateway prefix for the local riften-ipfs node, tried before public gateways. + /// Empty disables it. + riften_gateway: String, } const LOOP_SLEEP_TIME: Duration = Duration::from_secs(10); @@ -97,6 +135,7 @@ async fn fetch_bcmr( client: &reqwest::Client, urls: &[String], expected_hash: &str, + riften_gateway: &str, ) -> ( Option<(String, String)>, String, /* error */ @@ -106,9 +145,15 @@ async fn fetch_bcmr( let mut errors: Vec = Vec::new(); for url in urls { - for resolved in resolve_uri_candidates(url) { + for resolved in resolve_uri_candidates(url, riften_gateway) { + // The local gateway gets a short timeout so a hang falls back to public gateways fast. + let timeout = if !riften_gateway.is_empty() && resolved.starts_with(riften_gateway) { + LOCAL_GATEWAY_TIMEOUT + } else { + DOWNLOAD_TIMEOUT + }; let (contents, actual_hash) = - match get_url(client, &resolved, DOWNLOAD_TIMEOUT, MAX_BCMR_SIZE).await { + match get_url(client, &resolved, timeout, MAX_BCMR_SIZE).await { Ok(c) => c, Err(e) => { errors.push(e.to_string()); @@ -138,7 +183,12 @@ async fn fetch_bcmr( } } -async fn process_entry(client: &reqwest::Client, db: &SqlitePool, entry: &AuthChainEntry) { +async fn process_entry( + client: &reqwest::Client, + db: &SqlitePool, + entry: &AuthChainEntry, + riften_gateway: &str, +) { info!( "bcmr: Downloading BCMR for token {}, txid {}", entry.token_id, entry.txid @@ -163,7 +213,8 @@ async fn process_entry(client: &reqwest::Client, db: &SqlitePool, entry: &AuthCh } }; - let (json_str, error, is_fatal) = fetch_bcmr(client, &bcmr.uris, &hex::encode(bcmr.hash)).await; + let (json_str, error, is_fatal) = + fetch_bcmr(client, &bcmr.uris, &hex::encode(bcmr.hash), riften_gateway).await; let (json_str, actual_hash) = match json_str { Some(j) => j, @@ -234,17 +285,24 @@ async fn process_entry(client: &reqwest::Client, db: &SqlitePool, entry: &AuthCh } impl BCMRDownloader { - pub fn new(db: SqlitePool) -> Self { + pub fn new(db: SqlitePool, riften_gateway: String) -> Self { + if riften_gateway.is_empty() { + info!("bcmr: riften IPFS gateway not configured; using public gateways only"); + } else { + info!("bcmr: using riften IPFS gateway {riften_gateway} before public gateways"); + } Self { db, keep_running: Arc::new(AtomicBool::new(true)), download_task: None, + riften_gateway, } } pub fn start(&mut self) -> Result<()> { let db = self.db.clone(); let keep_running = self.keep_running.clone(); + let riften_gateway = self.riften_gateway.clone(); self.download_task = Some(tokio::spawn(async move { let client = reqwest::Client::new(); @@ -274,7 +332,7 @@ impl BCMRDownloader { for chunk in queue.chunks(MAX_PARALLEL_DOWNLOADS) { let futures: Vec<_> = chunk .iter() - .map(|entry| process_entry(&client, &db, entry)) + .map(|entry| process_entry(&client, &db, entry, &riften_gateway)) .collect(); futures::future::join_all(futures).await; } @@ -304,7 +362,7 @@ mod tests { #[test] fn bare_domain_resolves_to_well_known_then_root() { - let candidates = resolve_uri_candidates("chipnet-16.paryonusd.com"); + let candidates = resolve_uri_candidates("chipnet-16.paryonusd.com", ""); assert_eq!( candidates, vec![ @@ -317,7 +375,7 @@ mod tests { #[test] fn bare_domain_with_scheme_and_trailing_slash() { - let candidates = resolve_uri_candidates("https://example.com/"); + let candidates = resolve_uri_candidates("https://example.com/", ""); assert_eq!( candidates, vec![ @@ -329,7 +387,7 @@ mod tests { #[test] fn url_with_path_is_used_as_is() { - let candidates = resolve_uri_candidates("https://example.com/path/registry.json"); + let candidates = resolve_uri_candidates("https://example.com/path/registry.json", ""); assert_eq!( candidates, vec!["https://example.com/path/registry.json".to_string()] @@ -338,7 +396,7 @@ mod tests { #[test] fn schemeless_url_with_path_gets_https() { - let candidates = resolve_uri_candidates("example.com/registry.json"); + let candidates = resolve_uri_candidates("example.com/registry.json", ""); assert_eq!( candidates, vec!["https://example.com/registry.json".to_string()] @@ -346,16 +404,54 @@ mod tests { } #[test] - fn ipfs_uri_uses_a_gateway() { - let candidates = resolve_uri_candidates("ipfs://QmAbc123"); + fn ipfs_uri_uses_a_public_gateway_when_no_riften_gateway() { + let candidates = resolve_uri_candidates("ipfs://QmAbc123", ""); assert_eq!(candidates.len(), 1); assert!(candidates[0].ends_with("/QmAbc123")); assert!(IPFS_GATEWAYS.iter().any(|gw| candidates[0].starts_with(gw))); } + #[test] + fn ipfs_uri_tries_riften_gateway_first() { + let candidates = resolve_uri_candidates("ipfs://QmAbc123", "http://127.0.0.1:3002/ipfs/"); + assert_eq!(candidates.len(), 2); + assert_eq!(candidates[0], "http://127.0.0.1:3002/ipfs/QmAbc123"); + assert!(candidates[1].ends_with("/QmAbc123")); + assert!(IPFS_GATEWAYS.iter().any(|gw| candidates[1].starts_with(gw))); + } + + #[test] + fn ipfs_path_validation_accepts_real_cids_and_rejects_injection() { + // Real shapes seen on-chain: bare CIDv0/v1, directory + sub-path, and the libriften + // base64url CID with a filename extension appended directly. + assert!(is_safe_ipfs_path("QmAbc123")); + assert!(is_safe_ipfs_path("bafybeig4n5ut/icon.png")); + assert!(is_safe_ipfs_path( + "uAVUSIPZIwarHcHFcf3zv3HbVao50Xc8nByzxhCAPW2PMFUi5.png" + )); + // Traversal / URL-structure injection must be rejected. + assert!(!is_safe_ipfs_path("")); + assert!(!is_safe_ipfs_path("../admin")); + assert!(!is_safe_ipfs_path("cid/../../secret")); + assert!(!is_safe_ipfs_path("x@evil.com")); + assert!(!is_safe_ipfs_path("1.2.3.4:8080/x")); + assert!(!is_safe_ipfs_path("/etc/passwd")); + assert!(!is_safe_ipfs_path(".hidden")); + assert!(!is_safe_ipfs_path("cid/with space")); + assert!(!is_safe_ipfs_path("cid\r\nHost: evil")); + } + + #[test] + fn ipfs_uri_with_unsafe_path_yields_no_candidates() { + assert!( + resolve_uri_candidates("ipfs://../admin", "http://127.0.0.1:3002/ipfs/").is_empty() + ); + assert!(resolve_uri_candidates("ipfs://x@evil.com", "").is_empty()); + } + #[test] fn domain_with_port_is_treated_as_bare() { - let candidates = resolve_uri_candidates("example.com:8443"); + let candidates = resolve_uri_candidates("example.com:8443", ""); assert_eq!( candidates, vec![ diff --git a/src/db/bcmr/mod.rs b/src/db/bcmr/mod.rs index ab2dfd5..4462538 100644 --- a/src/db/bcmr/mod.rs +++ b/src/db/bcmr/mod.rs @@ -14,7 +14,6 @@ use sqlx::{Row, SqlitePool}; const MAX_DOWNLOAD_ATTEMPTS: usize = 100; -#[allow(dead_code)] pub struct AuthChainEntry { pub utxo: OutPointHash, pub txid: Txid, @@ -350,6 +349,35 @@ pub async fn update_bcmr_failure( Ok(()) } +/// Return the current authhead of a token's auth chain: the txid of the latest (highest +/// block height) auth_chain_entry. The authhead UTXO is always output 0 of that tx +/// (see `compute_outpoint_hash(&txid, 0)` in the indexer), so the outpoint is `:0`. +/// Returns None if the token has no auth chain entry. +pub async fn get_current_authhead( + pool: &SqlitePool, + token_hex: &str, +) -> Result> { + // txid is a deterministic tiebreaker: a forked/branched chain can have two entries at the + // same max height for one token, and without it SQLite's LIMIT 1 could flip between them. + let sql = "SELECT txid, height FROM auth_chain_entry + WHERE token_id = ? ORDER BY height DESC, txid DESC LIMIT 1"; + let token_blob = display_hex_to_blob::(token_hex)?; + + let row = sqlx::query(sql) + .bind(token_blob) + .fetch_optional(pool) + .await?; + + if let Some(r) = row { + let txid_blob: Vec = r.get(0); + let height: i64 = r.get(1); + let txid = Txid::from_blob(&txid_blob).context("failed to decode txid blob")?; + Ok(Some((txid, height as usize))) + } else { + Ok(None) + } +} + pub async fn get_token_bcmr(pool: &SqlitePool, token_hex: &str) -> Result> { let sql = r#" SELECT symbol, decimals, name, description, icon, web, actual_hash, expected_hash diff --git a/src/main.rs b/src/main.rs index 0245f82..29f76c5 100644 --- a/src/main.rs +++ b/src/main.rs @@ -346,7 +346,8 @@ async fn start_program( } }); - let mut bcmrdownloader = BCMRDownloader::new(db.bcmr_w.clone()); + let mut bcmrdownloader = + BCMRDownloader::new(db.bcmr_w.clone(), config.riften_ipfs_gateway.clone()); bcmrdownloader.start()?; let mut wellknowndownloader = WellKnownDownloader::new(db.bcmr_w.clone()); @@ -606,7 +607,11 @@ async fn launch() -> _ { ) .mount( "/bcmr", - routes![rpc::bcmr::token_bcmr, rpc::bcmr::token_bcmr_all], + routes![ + rpc::bcmr::token_bcmr, + rpc::bcmr::token_bcmr_all, + rpc::bcmr::token_authhead + ], ) .mount( "/oracle", diff --git a/src/rpc/bcmr.rs b/src/rpc/bcmr.rs index 1eda675..60f4e5b 100644 --- a/src/rpc/bcmr.rs +++ b/src/rpc/bcmr.rs @@ -10,9 +10,12 @@ use serde_json::json; use serde_json::Value; use crate::db::bcmr::get_well_known_bcmr; -use crate::db::{bcmr::get_token_bcmr, DB}; +use crate::db::{ + bcmr::{get_current_authhead, get_token_bcmr}, + DB, +}; use crate::rpc::err::{bad_request, db_error, ApiErrorCode, CachedApiResult}; -use crate::rpc::response::{cached_ok, CACHE_BCMR}; +use crate::rpc::response::{cached_ok, CACHE_BCMR, CACHE_NONE}; /// Fetches BCMR data for token from on-chain registry. /// Status: Stable @@ -93,3 +96,52 @@ pub async fn token_bcmr_all(category: Option<&str>, db: &State) -> CachedApi Ok(cached_ok(json!(bcmr_entries), CACHE_BCMR)) } + +/// Returns the current authhead of a token's auth chain — the outpoint that the next +/// BCMR-revision transaction must spend. The authhead is always output 0 of the latest +/// auth chain transaction, so `outpoint` is `:0`. +/// +/// Not cached: consumers (e.g. the riften-ipfs staging-pin flow) need the live head to +/// verify that an unbroadcast revision tx spends it. +/// +/// Returns `null` if the token has no auth chain entry. +/// +/// **Response Example:** +/// +/// ```json +/// { +/// "txid": "8289fa09272ac7930e219a2314965096a441f2f2fffdddd1f766f2cd3734fe28", +/// "vout": 0, +/// "outpoint": "8289fa09272ac7930e219a2314965096a441f2f2fffdddd1f766f2cd3734fe28:0", +/// "height": 123456 +/// } +/// ``` +#[get("/token//authhead")] +pub async fn token_authhead(category: Option<&str>, db: &State) -> CachedApiResult { + let token_id_hex = category + .context("category missing") + .map_err(|e| bad_request(ApiErrorCode::MissingCategory, &format!("Error: {e}")))?; + + let token_id = token_id_hex.parse::().map_err(|e| { + bad_request( + ApiErrorCode::InvalidTokenId, + &format!("Invalid token ID: {e}"), + ) + })?; + + let head = get_current_authhead(&db.bcmr_r, &token_id.to_string()) + .await + .map_err(db_error)?; + + let body = match head { + Some((txid, height)) => json!({ + "txid": txid.to_string(), + "vout": 0, + "outpoint": format!("{txid}:0"), + "height": height, + }), + None => Value::Null, + }; + + Ok(cached_ok(body, CACHE_NONE)) +}