wallet_endpoint.rs (42937B)
1 use std::{ 2 collections::{BTreeMap, BTreeSet}, 3 net::SocketAddr, 4 sync::Arc, 5 }; 6 7 use anyhow::{Context, Result}; 8 use axum::{ 9 Json, Router, 10 body::Body, 11 extract::{DefaultBodyLimit, Path, Query, State}, 12 http::{HeaderValue, Method, Request, StatusCode, header}, 13 middleware::{self, Next}, 14 response::{IntoResponse, Response}, 15 routing::{get, post}, 16 }; 17 use serde::{Deserialize, Serialize}; 18 use tokio::{net::TcpListener, sync::Semaphore}; 19 20 use crate::{ 21 adapters::{ 22 http::{ 23 types::{WalletTransactionContext, WalletTransactionFilters, WalletTransactionRow}, 24 ui::{ 25 add_pending_outputs, add_pending_v2_outputs, mark_wallet_reward_row, 26 transaction_v2_input_outpoints, wallet_transaction_row_for_addresses, 27 wallet_transaction_rows, wallet_transaction_v2_row, wallet_transaction_v2_rows, 28 }, 29 }, 30 p2p::GossipNetwork, 31 ui_data_store::SqliteUiDataStore, 32 }, 33 app::{NETWORK_ID, PROTOCOL_VERSION, SharedNode}, 34 domain::{ 35 AddressNetwork, AddressVersion, Amount, DEFAULT_FEE_PER_BYTE, 36 HYBRID_EXTERNAL_ADDRESS_GAP_LIMIT, MAX_BLOCK_BYTES, MICRO_IUNA, OutPoint, 37 TRANSACTION_SIGNING_FORMAT_VERSION, TRANSACTION_SIGNING_V1_ACTIVATION_HEIGHT, 38 TRANSACTION_V2_ACTIVATION_HEIGHT, TRANSACTION_V2_WIRE_VERSION, Transaction, 39 TransactionSubmitOutcome, TransactionV2, TxOutput, decode_hex, encode_versioned_address, 40 hex_encode, transaction_v2_is_active, 41 }, 42 }; 43 44 // A canonical v2 envelope is hexadecimal inside JSON, so its HTTP representation can be a little 45 // over twice the consensus block budget. 46 const MAX_TRANSACTION_BODY_BYTES: usize = MAX_BLOCK_BYTES * 2 + 4 * 1024; 47 const MAX_CONCURRENT_TRANSACTION_SUBMISSIONS: usize = 32; 48 // One legacy address plus the maximum external and reward address branches. 49 const MAX_WALLET_ADDRESSES: usize = 20_001; 50 51 #[derive(Clone)] 52 struct WalletEndpointState { 53 node: SharedNode, 54 gossip: GossipNetwork, 55 transaction_slots: Arc<Semaphore>, 56 ui_data_store: SqliteUiDataStore, 57 } 58 59 #[derive(Debug, Serialize)] 60 struct ApiError { 61 error: String, 62 } 63 64 #[derive(Debug, Serialize)] 65 struct WalletEndpointStatus { 66 api_version: u16, 67 ready: bool, 68 network_migration_required: bool, 69 network_id: String, 70 protocol_version: u32, 71 chain_id: String, 72 genesis_hash: String, 73 height: u64, 74 tip_hash: String, 75 transaction_signing_format_version: u16, 76 transaction_signing_v1_activation_height: u64, 77 default_fee_per_byte: Amount, 78 amount_unit: &'static str, 79 microiuna_per_iuna: Amount, 80 transaction_v2_wire_version: u16, 81 transaction_v2_activation_height: Option<u64>, 82 transaction_v2_active: bool, 83 hybrid_address_gap_limit: u32, 84 } 85 86 #[derive(Debug, Deserialize)] 87 struct WalletSnapshotRequest { 88 addresses: Vec<String>, 89 } 90 91 #[derive(Debug, Serialize)] 92 struct WalletSnapshotResponse { 93 addresses: Vec<WalletAddressSnapshot>, 94 confirmed: Amount, 95 spendable: Amount, 96 pending_outgoing: Amount, 97 pending_incoming: Amount, 98 height: u64, 99 tip_hash: String, 100 utxos: Vec<WalletSnapshotUtxo>, 101 } 102 103 #[derive(Debug, Serialize)] 104 struct WalletAddressSnapshot { 105 address: String, 106 version: u8, 107 used: bool, 108 confirmed: Amount, 109 spendable: Amount, 110 } 111 112 #[derive(Debug, Serialize)] 113 struct WalletSnapshotUtxo { 114 address: String, 115 outpoint: OutPoint, 116 output: TxOutput, 117 } 118 119 #[derive(Debug, Deserialize)] 120 struct WalletTransactionsRequest { 121 addresses: Vec<String>, 122 #[serde(default = "default_true")] 123 transfer: bool, 124 #[serde(default = "default_true")] 125 mine: bool, 126 #[serde(default = "default_true")] 127 burn: bool, 128 #[serde(default = "default_true")] 129 reward: bool, 130 #[serde(default)] 131 offset: usize, 132 limit: Option<usize>, 133 } 134 135 #[derive(Debug, Serialize)] 136 struct WalletTransactionsResponse { 137 items: Vec<WalletTransactionRow>, 138 offset: usize, 139 limit: usize, 140 total: usize, 141 has_more: bool, 142 next_offset: Option<usize>, 143 } 144 145 #[derive(Debug, Serialize)] 146 struct BalanceResponse { 147 address: String, 148 public_key: String, 149 confirmed: Amount, 150 spendable: Amount, 151 pending_outgoing: Amount, 152 pending_incoming: Amount, 153 height: u64, 154 tip_hash: String, 155 } 156 157 #[derive(Debug, Serialize)] 158 struct UtxoResponse { 159 address: String, 160 public_key: String, 161 height: u64, 162 tip_hash: String, 163 utxos: Vec<WalletUtxo>, 164 } 165 166 #[derive(Debug, Serialize)] 167 struct WalletUtxo { 168 outpoint: OutPoint, 169 output: TxOutput, 170 } 171 172 #[derive(Debug, Serialize)] 173 struct SubmitResponse { 174 transaction_id: String, 175 status: &'static str, 176 } 177 178 #[derive(Debug, Deserialize)] 179 struct SubmitV2Request { 180 envelope: String, 181 } 182 183 #[derive(Debug, Serialize)] 184 struct TransactionResponse { 185 transaction_id: String, 186 status: &'static str, 187 transaction: Transaction, 188 } 189 190 #[derive(Debug, Default, Deserialize)] 191 struct TransactionsQuery { 192 tx: Option<bool>, 193 mine: Option<bool>, 194 burn: Option<bool>, 195 reward: Option<bool>, 196 offset: Option<usize>, 197 limit: Option<usize>, 198 } 199 200 #[derive(Clone, Copy)] 201 struct TransactionFilters { 202 transfer: bool, 203 mine: bool, 204 burn: bool, 205 reward: bool, 206 } 207 208 impl TransactionFilters { 209 fn from_query(query: &TransactionsQuery) -> Self { 210 Self { 211 transfer: query.tx.unwrap_or(true), 212 // Keep the public endpoint's unfiltered default backwards compatible. 213 // The wallet sends all four choices explicitly. 214 mine: query.mine.unwrap_or(true), 215 burn: query.burn.unwrap_or(true), 216 reward: query.reward.unwrap_or(true), 217 } 218 } 219 220 fn allows(self, transaction: &Transaction) -> bool { 221 match transaction { 222 Transaction::Transfer { .. } => self.transfer, 223 Transaction::Mine { .. } => self.mine, 224 Transaction::Burn { .. } => self.burn, 225 } 226 } 227 228 fn kinds(self) -> Vec<&'static str> { 229 let mut kinds = Vec::with_capacity(4); 230 if self.transfer { 231 kinds.push("transfer"); 232 } 233 if self.mine { 234 kinds.push("mine"); 235 } 236 if self.burn { 237 kinds.push("burn"); 238 } 239 if self.reward { 240 kinds.push("reward"); 241 } 242 kinds 243 } 244 } 245 246 #[derive(Debug, Serialize)] 247 struct TransactionsResponse { 248 items: Vec<AddressTransaction>, 249 offset: usize, 250 limit: usize, 251 total: usize, 252 has_more: bool, 253 next_offset: Option<usize>, 254 } 255 256 #[derive(Debug, Serialize)] 257 struct AddressTransaction { 258 kind: String, 259 status: &'static str, 260 block_height: Option<u64>, 261 timestamp_ms: Option<u64>, 262 transaction: Transaction, 263 } 264 265 pub async fn serve( 266 node: SharedNode, 267 gossip: GossipNetwork, 268 ui_data_store: SqliteUiDataStore, 269 addr: SocketAddr, 270 ) -> Result<()> { 271 let listener = TcpListener::bind(addr) 272 .await 273 .with_context(|| format!("binding public wallet endpoint on {addr}"))?; 274 println!("wallet endpoint: http://{}", listener.local_addr()?); 275 axum::serve(listener, router(node, gossip, ui_data_store)) 276 .await 277 .context("serving public wallet endpoint") 278 } 279 280 fn router(node: SharedNode, gossip: GossipNetwork, ui_data_store: SqliteUiDataStore) -> Router { 281 Router::new() 282 .route("/v1/status", get(status)) 283 .route("/v1/wallets/snapshot", post(wallet_snapshot)) 284 .route("/v1/wallets/transactions", post(wallet_transactions)) 285 .route("/v1/addresses/{address}/balance", get(balance)) 286 .route("/v1/addresses/{address}/utxos", get(utxos)) 287 .route( 288 "/v1/addresses/{address}/transactions", 289 get(address_transactions), 290 ) 291 .route("/v1/transactions/{transaction_id}", get(transaction)) 292 .route("/v1/transactions", post(submit_transaction)) 293 .route("/v1/transactions-v2", post(submit_transaction_v2)) 294 .layer(DefaultBodyLimit::max(MAX_TRANSACTION_BODY_BYTES)) 295 .layer(middleware::from_fn(cors)) 296 .with_state(WalletEndpointState { 297 node, 298 gossip, 299 transaction_slots: Arc::new(Semaphore::new(MAX_CONCURRENT_TRANSACTION_SUBMISSIONS)), 300 ui_data_store, 301 }) 302 } 303 304 async fn status(State(state): State<WalletEndpointState>) -> Json<WalletEndpointStatus> { 305 let node = state.node.lock().await; 306 let ledger = node.ledger(); 307 Json(WalletEndpointStatus { 308 api_version: 2, 309 ready: node.has_real_chain() && node.network_migration_from().is_none(), 310 network_migration_required: node.network_migration_from().is_some(), 311 network_id: NETWORK_ID.to_string(), 312 protocol_version: PROTOCOL_VERSION, 313 chain_id: ledger.launch_profile().profile_id.clone(), 314 genesis_hash: ledger.genesis_hash().to_string(), 315 height: ledger.height(), 316 tip_hash: ledger.tip_hash().to_string(), 317 transaction_signing_format_version: TRANSACTION_SIGNING_FORMAT_VERSION, 318 transaction_signing_v1_activation_height: TRANSACTION_SIGNING_V1_ACTIVATION_HEIGHT, 319 default_fee_per_byte: DEFAULT_FEE_PER_BYTE, 320 amount_unit: "microiuna", 321 microiuna_per_iuna: MICRO_IUNA, 322 transaction_v2_wire_version: TRANSACTION_V2_WIRE_VERSION, 323 transaction_v2_activation_height: TRANSACTION_V2_ACTIVATION_HEIGHT, 324 transaction_v2_active: transaction_v2_is_active(ledger.height().saturating_add(1)), 325 hybrid_address_gap_limit: HYBRID_EXTERNAL_ADDRESS_GAP_LIMIT, 326 }) 327 } 328 329 async fn wallet_snapshot( 330 State(state): State<WalletEndpointState>, 331 Json(request): Json<WalletSnapshotRequest>, 332 ) -> Result<Json<WalletSnapshotResponse>, (StatusCode, Json<ApiError>)> { 333 if request.addresses.is_empty() { 334 return Err(api_error( 335 StatusCode::BAD_REQUEST, 336 "at least one address is required", 337 )); 338 } 339 if request.addresses.len() > MAX_WALLET_ADDRESSES { 340 return Err(api_error( 341 StatusCode::BAD_REQUEST, 342 format!("a wallet snapshot accepts at most {MAX_WALLET_ADDRESSES} addresses"), 343 )); 344 } 345 346 let node = state.node.lock().await; 347 let ledger = node.ledger(); 348 let network = 349 crate::domain::AddressNetwork::from_profile_id(&ledger.launch_profile().profile_id); 350 let used_hybrid = ledger 351 .used_hybrid_encoded_addresses() 352 .map_err(internal_error)? 353 .into_iter() 354 .collect::<BTreeSet<_>>(); 355 let mut seen = BTreeSet::new(); 356 let mut requested_addresses = Vec::new(); 357 for requested in request.addresses { 358 let decoded = node.decode_user_address(&requested).map_err(bad_request)?; 359 let display = encode_versioned_address(decoded, network).map_err(bad_request)?; 360 if !seen.insert(display.clone()) { 361 continue; 362 } 363 let ledger_address = match decoded.version { 364 AddressVersion::Ed25519PublicKey => hex_encode(decoded.payload), 365 AddressVersion::HybridKeyCommitment => display.clone(), 366 }; 367 requested_addresses.push((display, decoded.version, ledger_address)); 368 } 369 let ledger_addresses = requested_addresses 370 .iter() 371 .map(|(_, _, address)| address.clone()) 372 .collect::<BTreeSet<_>>(); 373 let display_by_ledger_address = requested_addresses 374 .iter() 375 .map(|(display, _, ledger_address)| (ledger_address.clone(), display.clone())) 376 .collect::<BTreeMap<_, _>>(); 377 let mut confirmed_by_address = BTreeMap::<String, Amount>::new(); 378 for (_, output) in ledger.utxos_for_addresses(&ledger_addresses) { 379 let balance = confirmed_by_address.entry(output.address).or_default(); 380 *balance = balance.checked_add(output.amount).ok_or_else(|| { 381 api_error( 382 StatusCode::INTERNAL_SERVER_ERROR, 383 "wallet balance overflows", 384 ) 385 })?; 386 } 387 let available = ledger 388 .available_utxos_for_addresses(&ledger_addresses) 389 .map_err(internal_error)?; 390 let mut spendable_by_address = BTreeMap::<String, Amount>::new(); 391 for (_, output) in &available { 392 let balance = spendable_by_address 393 .entry(output.address.clone()) 394 .or_default(); 395 *balance = balance.checked_add(output.amount).ok_or_else(|| { 396 api_error( 397 StatusCode::INTERNAL_SERVER_ERROR, 398 "wallet balance overflows", 399 ) 400 })?; 401 } 402 let mut confirmed = 0_u64; 403 let mut spendable = 0_u64; 404 let addresses = requested_addresses 405 .into_iter() 406 .map(|(display, version, ledger_address)| { 407 let address_confirmed = confirmed_by_address 408 .get(&ledger_address) 409 .copied() 410 .unwrap_or_default(); 411 let address_spendable = spendable_by_address 412 .get(&ledger_address) 413 .copied() 414 .unwrap_or_default(); 415 confirmed = confirmed.checked_add(address_confirmed).ok_or_else(|| { 416 api_error( 417 StatusCode::INTERNAL_SERVER_ERROR, 418 "wallet balance overflows", 419 ) 420 })?; 421 spendable = spendable.checked_add(address_spendable).ok_or_else(|| { 422 api_error( 423 StatusCode::INTERNAL_SERVER_ERROR, 424 "wallet balance overflows", 425 ) 426 })?; 427 Ok(WalletAddressSnapshot { 428 address: display.clone(), 429 version: version.wire_id(), 430 used: version == AddressVersion::HybridKeyCommitment 431 && used_hybrid.contains(&display), 432 confirmed: address_confirmed, 433 spendable: address_spendable, 434 }) 435 }) 436 .collect::<Result<Vec<_>, (StatusCode, Json<ApiError>)>>()?; 437 let mut utxos = available 438 .into_iter() 439 .filter_map(|(outpoint, output)| { 440 display_by_ledger_address 441 .get(&output.address) 442 .cloned() 443 .map(|address| WalletSnapshotUtxo { 444 address, 445 outpoint, 446 output, 447 }) 448 }) 449 .collect::<Vec<_>>(); 450 utxos.sort_by(|left, right| { 451 right 452 .output 453 .amount 454 .cmp(&left.output.amount) 455 .then_with(|| left.outpoint.cmp(&right.outpoint)) 456 }); 457 Ok(Json(WalletSnapshotResponse { 458 addresses, 459 confirmed, 460 spendable, 461 pending_outgoing: confirmed.saturating_sub(spendable), 462 pending_incoming: spendable.saturating_sub(confirmed), 463 height: ledger.height(), 464 tip_hash: ledger.tip_hash().to_string(), 465 utxos, 466 })) 467 } 468 469 async fn wallet_transactions( 470 State(state): State<WalletEndpointState>, 471 Json(request): Json<WalletTransactionsRequest>, 472 ) -> Result<Json<WalletTransactionsResponse>, (StatusCode, Json<ApiError>)> { 473 if request.addresses.is_empty() || request.addresses.len() > MAX_WALLET_ADDRESSES { 474 return Err(api_error( 475 StatusCode::BAD_REQUEST, 476 format!("provide between 1 and {MAX_WALLET_ADDRESSES} wallet addresses"), 477 )); 478 } 479 let offset = request.offset; 480 let limit = request.limit.unwrap_or(25).clamp(1, 100); 481 let filters = WalletTransactionFilters { 482 transfer: request.transfer, 483 mine: request.mine, 484 burn: request.burn, 485 reward: request.reward, 486 }; 487 let kinds = wallet_transaction_kinds(filters); 488 let (addresses, pending, pending_v2, domain, network, mut pending_outputs) = { 489 let node = state.node.lock().await; 490 let ledger = node.ledger(); 491 let network = AddressNetwork::from_profile_id(&ledger.launch_profile().profile_id); 492 let mut addresses = Vec::new(); 493 for requested in request.addresses { 494 let decoded = node.decode_user_address(&requested).map_err(bad_request)?; 495 let address = match decoded.version { 496 AddressVersion::Ed25519PublicKey => hex_encode(decoded.payload), 497 AddressVersion::HybridKeyCommitment => { 498 encode_versioned_address(decoded, network).map_err(bad_request)? 499 } 500 }; 501 if !addresses.contains(&address) { 502 addresses.push(address); 503 } 504 } 505 let pending_v2 = node.pending_transactions_v2(); 506 let pending_outputs = pending_v2 507 .iter() 508 .flat_map(transaction_v2_input_outpoints) 509 .filter_map(|outpoint| { 510 ledger 511 .output_for_outpoint(&outpoint) 512 .map(|output| (outpoint, output)) 513 }) 514 .collect::<BTreeMap<_, _>>(); 515 ( 516 addresses, 517 node.pending_transactions(), 518 pending_v2, 519 ledger.transaction_v2_domain().map_err(internal_error)?, 520 network, 521 pending_outputs, 522 ) 523 }; 524 525 let mut pending_outpoints = BTreeSet::new(); 526 collect_legacy_input_outpoints(pending.iter(), &mut pending_outpoints); 527 pending_outpoints.extend(pending_v2.iter().flat_map(transaction_v2_input_outpoints)); 528 if !pending_outpoints.is_empty() { 529 let store = state.ui_data_store.clone(); 530 let stored = tokio::task::spawn_blocking(move || store.load_outputs(&pending_outpoints)) 531 .await 532 .map_err(internal_error)? 533 .map_err(internal_error)?; 534 pending_outputs.extend(stored); 535 } 536 add_pending_outputs(&mut pending_outputs, &pending); 537 add_pending_v2_outputs(&mut pending_outputs, &pending_v2, &domain, network) 538 .map_err(internal_error)?; 539 let mut pending_rows = wallet_transaction_v2_rows( 540 &addresses, 541 &pending_v2, 542 &pending_outputs, 543 filters, 544 &domain, 545 network, 546 ); 547 pending_rows.extend(wallet_transaction_rows( 548 &addresses, 549 pending, 550 &[], 551 &pending_outputs, 552 filters, 553 )); 554 let pending_total = pending_rows.len(); 555 let mut items = pending_rows 556 .into_iter() 557 .skip(offset.min(pending_total)) 558 .take(limit) 559 .collect::<Vec<_>>(); 560 561 let confirmed_offset = offset.saturating_sub(pending_total); 562 let remaining = limit.saturating_sub(items.len()); 563 let fetch_limit = confirmed_offset.saturating_add(remaining); 564 let store = state.ui_data_store.clone(); 565 let query_addresses = addresses.clone(); 566 let ((legacy_rows, legacy_total), (v2_rows, v2_total)) = 567 tokio::task::spawn_blocking(move || -> Result<_> { 568 Ok(( 569 store.load_wallet_transactions_for_addresses( 570 &query_addresses, 571 &kinds, 572 0, 573 fetch_limit, 574 )?, 575 store.load_wallet_transactions_v2(&query_addresses, &kinds, 0, fetch_limit)?, 576 )) 577 }) 578 .await 579 .map_err(internal_error)? 580 .map_err(internal_error)?; 581 582 let decoded_v2 = v2_rows 583 .into_iter() 584 .filter_map(|row| { 585 let bytes = decode_hex(&row.envelope).ok()?; 586 let (decoded_domain, transaction) = TransactionV2::decode(&bytes).ok()?; 587 (decoded_domain == domain).then_some((row, transaction)) 588 }) 589 .collect::<Vec<_>>(); 590 let mut confirmed_outpoints = BTreeSet::new(); 591 collect_legacy_input_outpoints( 592 legacy_rows.iter().map(|row| &row.transaction), 593 &mut confirmed_outpoints, 594 ); 595 confirmed_outpoints.extend( 596 decoded_v2 597 .iter() 598 .flat_map(|(_, transaction)| transaction_v2_input_outpoints(transaction)), 599 ); 600 let confirmed_outputs = if confirmed_outpoints.is_empty() { 601 BTreeMap::new() 602 } else { 603 let store = state.ui_data_store.clone(); 604 tokio::task::spawn_blocking(move || store.load_outputs(&confirmed_outpoints)) 605 .await 606 .map_err(internal_error)? 607 .map_err(internal_error)? 608 }; 609 let reward_blocks = { 610 let node = state.node.lock().await; 611 legacy_rows 612 .iter() 613 .filter(|row| row.kind == "reward") 614 .filter_map(|row| { 615 let block = node.chain().get(row.block_height as usize)?.clone(); 616 Some((row.block_height, block)) 617 }) 618 .collect::<BTreeMap<_, _>>() 619 }; 620 let mut confirmed = legacy_rows 621 .into_iter() 622 .filter_map(|row| { 623 let is_reward = row.kind == "reward"; 624 let mut item = wallet_transaction_row_for_addresses( 625 &addresses, 626 &row.transaction, 627 &confirmed_outputs, 628 &WalletTransactionContext { 629 status: "confirmed", 630 block_height: Some(row.block_height), 631 timestamp_ms: Some(row.timestamp_ms), 632 block_finalizer: Some(row.block_finalizer), 633 }, 634 )?; 635 if is_reward { 636 mark_wallet_reward_row(&mut item, reward_blocks.get(&row.block_height)); 637 } 638 Some((row.sort_key, item)) 639 }) 640 .collect::<Vec<_>>(); 641 confirmed.extend(decoded_v2.into_iter().filter_map(|(row, transaction)| { 642 wallet_transaction_v2_row( 643 &addresses, 644 &transaction, 645 &confirmed_outputs, 646 &domain, 647 network, 648 &WalletTransactionContext { 649 status: "confirmed", 650 block_height: Some(row.block_height), 651 timestamp_ms: Some(row.timestamp_ms), 652 block_finalizer: Some(row.block_finalizer), 653 }, 654 ) 655 .ok() 656 .flatten() 657 .map(|item| (row.sort_key, item)) 658 })); 659 confirmed.sort_by(|left, right| right.0.cmp(&left.0)); 660 items.extend( 661 confirmed 662 .into_iter() 663 .skip(confirmed_offset) 664 .take(remaining) 665 .map(|(_, item)| item), 666 ); 667 let total = pending_total 668 .saturating_add(legacy_total) 669 .saturating_add(v2_total); 670 let next_offset = offset.saturating_add(items.len()); 671 Ok(Json(WalletTransactionsResponse { 672 items, 673 offset: offset.min(total), 674 limit, 675 total, 676 has_more: next_offset < total, 677 next_offset: (next_offset < total).then_some(next_offset), 678 })) 679 } 680 681 fn default_true() -> bool { 682 true 683 } 684 685 fn wallet_transaction_kinds(filters: WalletTransactionFilters) -> Vec<&'static str> { 686 let mut kinds = Vec::new(); 687 if filters.transfer { 688 kinds.push("transfer"); 689 } 690 if filters.mine { 691 kinds.push("mine"); 692 } 693 if filters.burn { 694 kinds.push("burn"); 695 } 696 if filters.reward { 697 kinds.push("reward"); 698 } 699 kinds 700 } 701 702 fn collect_legacy_input_outpoints<'a>( 703 transactions: impl IntoIterator<Item = &'a Transaction>, 704 outpoints: &mut BTreeSet<OutPoint>, 705 ) { 706 for transaction in transactions { 707 if let Transaction::Transfer { inputs, .. } | Transaction::Burn { inputs, .. } = transaction 708 { 709 outpoints.extend(inputs.iter().map(|input| input.outpoint.clone())); 710 } 711 } 712 } 713 714 async fn balance( 715 State(state): State<WalletEndpointState>, 716 Path(address): Path<String>, 717 ) -> Result<Json<BalanceResponse>, (StatusCode, Json<ApiError>)> { 718 let node = state.node.lock().await; 719 let public_key = node.normalize_user_address(&address).map_err(bad_request)?; 720 let ledger = node.ledger(); 721 let confirmed = ledger.balance_of(&public_key); 722 let spendable = ledger 723 .available_utxos_for_address(&public_key) 724 .map_err(internal_error)? 725 .iter() 726 .map(|(_, output)| output.amount) 727 .sum(); 728 Ok(Json(BalanceResponse { 729 address: address.trim().to_ascii_lowercase(), 730 public_key, 731 confirmed, 732 spendable, 733 pending_outgoing: confirmed.saturating_sub(spendable), 734 pending_incoming: spendable.saturating_sub(confirmed), 735 height: ledger.height(), 736 tip_hash: ledger.tip_hash().to_string(), 737 })) 738 } 739 740 async fn utxos( 741 State(state): State<WalletEndpointState>, 742 Path(address): Path<String>, 743 ) -> Result<Json<UtxoResponse>, (StatusCode, Json<ApiError>)> { 744 let node = state.node.lock().await; 745 let public_key = node.normalize_user_address(&address).map_err(bad_request)?; 746 let ledger = node.ledger(); 747 let utxos = ledger 748 .available_utxos_for_address(&public_key) 749 .map_err(internal_error)? 750 .into_iter() 751 .map(|(outpoint, output)| WalletUtxo { outpoint, output }) 752 .collect(); 753 Ok(Json(UtxoResponse { 754 address: address.trim().to_ascii_lowercase(), 755 public_key, 756 height: ledger.height(), 757 tip_hash: ledger.tip_hash().to_string(), 758 utxos, 759 })) 760 } 761 762 async fn transaction( 763 State(state): State<WalletEndpointState>, 764 Path(transaction_id): Path<String>, 765 ) -> Result<Json<TransactionResponse>, (StatusCode, Json<ApiError>)> { 766 let pending = { 767 let node = state.node.lock().await; 768 let ledger = node.ledger(); 769 ledger 770 .pending() 771 .iter() 772 .find(|item| item.signature() == transaction_id) 773 .cloned() 774 .map(|transaction| ("pending", transaction)) 775 .or_else(|| { 776 ledger 777 .orphan_transactions() 778 .iter() 779 .find(|item| item.signature() == transaction_id) 780 .cloned() 781 .map(|transaction| ("orphan", transaction)) 782 }) 783 }; 784 let (status, transaction) = match pending { 785 Some(transaction) => transaction, 786 None => { 787 let store = state.ui_data_store.clone(); 788 let lookup_id = transaction_id.clone(); 789 let transaction = tokio::task::spawn_blocking(move || { 790 store.load_wallet_transaction_by_signature(&lookup_id) 791 }) 792 .await 793 .map_err(internal_error)? 794 .map_err(internal_error)? 795 .ok_or_else(|| api_error(StatusCode::NOT_FOUND, "transaction not found"))?; 796 ("confirmed", transaction.transaction) 797 } 798 }; 799 Ok(Json(TransactionResponse { 800 transaction_id, 801 status, 802 transaction, 803 })) 804 } 805 806 async fn address_transactions( 807 State(state): State<WalletEndpointState>, 808 Path(address): Path<String>, 809 Query(query): Query<TransactionsQuery>, 810 ) -> Result<Json<TransactionsResponse>, (StatusCode, Json<ApiError>)> { 811 let offset = query.offset.unwrap_or(0); 812 let limit = query.limit.unwrap_or(25).clamp(1, 100); 813 let filters = TransactionFilters::from_query(&query); 814 let (public_key, pending) = { 815 let node = state.node.lock().await; 816 let public_key = node.normalize_user_address(&address).map_err(bad_request)?; 817 let pending = node 818 .pending_transactions() 819 .into_iter() 820 .rev() 821 .filter(|transaction| transaction_mentions_address(transaction, &public_key)) 822 .filter(|transaction| filters.allows(transaction)) 823 .collect::<Vec<_>>(); 824 (public_key, pending) 825 }; 826 let pending_total = pending.len(); 827 let mut items = pending 828 .into_iter() 829 .skip(offset.min(pending_total)) 830 .take(limit) 831 .map(|transaction| AddressTransaction { 832 kind: transaction_kind(&transaction).to_string(), 833 status: "pending", 834 block_height: None, 835 timestamp_ms: None, 836 transaction, 837 }) 838 .collect::<Vec<_>>(); 839 840 let confirmed_offset = offset.saturating_sub(pending_total); 841 let remaining = limit.saturating_sub(items.len()); 842 let kinds = filters.kinds(); 843 let confirmed = if kinds.is_empty() { 844 (Vec::new(), 0) 845 } else { 846 let store = state.ui_data_store.clone(); 847 tokio::task::spawn_blocking(move || { 848 store.load_wallet_transactions(&public_key, &kinds, confirmed_offset, remaining) 849 }) 850 .await 851 .map_err(internal_error)? 852 .map_err(internal_error)? 853 }; 854 let confirmed_total = confirmed.1; 855 items.extend(confirmed.0.into_iter().map(|row| AddressTransaction { 856 kind: row.kind, 857 status: "confirmed", 858 block_height: Some(row.block_height), 859 timestamp_ms: Some(row.timestamp_ms), 860 transaction: row.transaction, 861 })); 862 let total = pending_total + confirmed_total; 863 let next_offset = offset.saturating_add(items.len()); 864 Ok(Json(TransactionsResponse { 865 items, 866 offset: offset.min(total), 867 limit, 868 total, 869 has_more: next_offset < total, 870 next_offset: (next_offset < total).then_some(next_offset), 871 })) 872 } 873 874 fn transaction_mentions_address(transaction: &Transaction, address: &str) -> bool { 875 match transaction { 876 Transaction::Transfer { 877 inputs, outputs, .. 878 } => { 879 inputs.iter().any(|input| input.owner == address) 880 || outputs.iter().any(|output| output.address == address) 881 } 882 Transaction::Burn { inputs, change, .. } => { 883 inputs.iter().any(|input| input.owner == address) 884 || change.iter().any(|output| output.address == address) 885 } 886 Transaction::Mine { recipient, .. } => recipient == address, 887 } 888 } 889 890 fn transaction_kind(transaction: &Transaction) -> &'static str { 891 match transaction { 892 Transaction::Transfer { .. } => "transfer", 893 Transaction::Burn { .. } => "burn", 894 Transaction::Mine { .. } => "mine", 895 } 896 } 897 898 async fn submit_transaction( 899 State(state): State<WalletEndpointState>, 900 Json(transaction): Json<Transaction>, 901 ) -> Result<(StatusCode, Json<SubmitResponse>), (StatusCode, Json<ApiError>)> { 902 let _permit = state 903 .transaction_slots 904 .clone() 905 .try_acquire_owned() 906 .map_err(|_| { 907 api_error( 908 StatusCode::TOO_MANY_REQUESTS, 909 "too many transaction submissions", 910 ) 911 })?; 912 let transaction_id = transaction.signature().to_string(); 913 let (outcome, outbox) = { 914 let mut node = state.node.lock().await; 915 if !node.has_real_chain() { 916 return Err(api_error( 917 StatusCode::SERVICE_UNAVAILABLE, 918 "node has not joined a chain yet", 919 )); 920 } 921 if node.network_migration_from().is_some() { 922 return Err(api_error( 923 StatusCode::SERVICE_UNAVAILABLE, 924 "node requires a network migration reset", 925 )); 926 } 927 let outcome = node 928 .submit_external_wallet_transaction(transaction) 929 .map_err(bad_request)?; 930 (outcome, node.drain_outbox()) 931 }; 932 if outcome.added() { 933 state 934 .gossip 935 .broadcast(outbox) 936 .await 937 .map_err(internal_error)?; 938 } 939 let (code, status) = match outcome { 940 TransactionSubmitOutcome::Added => (StatusCode::ACCEPTED, "accepted"), 941 TransactionSubmitOutcome::AlreadyKnown => (StatusCode::OK, "already_known"), 942 TransactionSubmitOutcome::ConflictsWithPending => (StatusCode::CONFLICT, "conflict"), 943 }; 944 Ok(( 945 code, 946 Json(SubmitResponse { 947 transaction_id, 948 status, 949 }), 950 )) 951 } 952 953 async fn submit_transaction_v2( 954 State(state): State<WalletEndpointState>, 955 Json(request): Json<SubmitV2Request>, 956 ) -> Result<(StatusCode, Json<SubmitResponse>), (StatusCode, Json<ApiError>)> { 957 let _permit = state 958 .transaction_slots 959 .clone() 960 .try_acquire_owned() 961 .map_err(|_| { 962 api_error( 963 StatusCode::TOO_MANY_REQUESTS, 964 "too many transaction submissions", 965 ) 966 })?; 967 let (transaction_id, outcome, outbox) = { 968 let mut node = state.node.lock().await; 969 if !node.has_real_chain() { 970 return Err(api_error( 971 StatusCode::SERVICE_UNAVAILABLE, 972 "node has not joined a chain yet", 973 )); 974 } 975 if node.network_migration_from().is_some() { 976 return Err(api_error( 977 StatusCode::SERVICE_UNAVAILABLE, 978 "node requires a network migration reset", 979 )); 980 } 981 let (transaction_id, outcome) = node 982 .submit_external_wallet_transaction_v2(&request.envelope) 983 .map_err(bad_request)?; 984 (transaction_id, outcome, node.drain_outbox()) 985 }; 986 if outcome.added() { 987 state 988 .gossip 989 .broadcast(outbox) 990 .await 991 .map_err(internal_error)?; 992 } 993 let (code, status) = match outcome { 994 TransactionSubmitOutcome::Added => (StatusCode::ACCEPTED, "accepted"), 995 TransactionSubmitOutcome::AlreadyKnown => (StatusCode::OK, "already_known"), 996 TransactionSubmitOutcome::ConflictsWithPending => (StatusCode::CONFLICT, "conflict"), 997 }; 998 Ok(( 999 code, 1000 Json(SubmitResponse { 1001 transaction_id, 1002 status, 1003 }), 1004 )) 1005 } 1006 1007 fn bad_request(error: impl std::fmt::Display) -> (StatusCode, Json<ApiError>) { 1008 api_error(StatusCode::BAD_REQUEST, error) 1009 } 1010 1011 fn internal_error(error: impl std::fmt::Display) -> (StatusCode, Json<ApiError>) { 1012 api_error(StatusCode::INTERNAL_SERVER_ERROR, error) 1013 } 1014 1015 fn api_error(status: StatusCode, error: impl std::fmt::Display) -> (StatusCode, Json<ApiError>) { 1016 ( 1017 status, 1018 Json(ApiError { 1019 error: error.to_string(), 1020 }), 1021 ) 1022 } 1023 1024 async fn cors(request: Request<Body>, next: Next) -> Response { 1025 if request.method() == Method::OPTIONS { 1026 return cors_headers(StatusCode::NO_CONTENT.into_response()); 1027 } 1028 cors_headers(next.run(request).await) 1029 } 1030 1031 fn cors_headers(mut response: Response) -> Response { 1032 let headers = response.headers_mut(); 1033 headers.insert( 1034 header::ACCESS_CONTROL_ALLOW_ORIGIN, 1035 HeaderValue::from_static("*"), 1036 ); 1037 headers.insert( 1038 header::ACCESS_CONTROL_ALLOW_METHODS, 1039 HeaderValue::from_static("GET, POST, OPTIONS"), 1040 ); 1041 headers.insert( 1042 header::ACCESS_CONTROL_ALLOW_HEADERS, 1043 HeaderValue::from_static("content-type"), 1044 ); 1045 response 1046 } 1047 1048 #[cfg(test)] 1049 mod tests { 1050 use std::{ 1051 collections::BTreeMap, 1052 net::{IpAddr, Ipv4Addr}, 1053 sync::Arc, 1054 }; 1055 1056 use axum::{body::to_bytes, http::Request}; 1057 use tokio::sync::Mutex; 1058 use tower::ServiceExt; 1059 1060 use crate::{ 1061 adapters::p2p::GossipNetwork, 1062 app::{NodeCore, PeerBook}, 1063 domain::{AddressNetwork, Ledger, MICRO_IUNA, Wallet, encode_address}, 1064 }; 1065 1066 use super::*; 1067 1068 async fn test_app() -> (Router, Wallet, Ledger, tempfile::TempDir) { 1069 let wallet = Wallet::from_seed("public-wallet-endpoint-test"); 1070 let ledger = Ledger::new( 1071 BTreeMap::from([(wallet.address().to_string(), 5 * MICRO_IUNA)]), 1072 1, 1073 ); 1074 let node = Arc::new(Mutex::new(NodeCore::from_ledger( 1075 wallet.clone(), 1076 ledger.clone(), 1077 0, 1078 ))); 1079 let peers = Arc::new(Mutex::new(PeerBook::default())); 1080 let gossip = GossipNetwork::start( 1081 node.clone(), 1082 peers, 1083 SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0), 1084 None, 1085 false, 1086 ) 1087 .await 1088 .unwrap(); 1089 let dir = tempfile::tempdir().unwrap(); 1090 let ui_data_store = SqliteUiDataStore::open(dir.path().join("ui.sqlite3")).unwrap(); 1091 ui_data_store 1092 .project_snapshot(&ledger.snapshot(), false) 1093 .unwrap(); 1094 (router(node, gossip, ui_data_store), wallet, ledger, dir) 1095 } 1096 1097 #[tokio::test] 1098 async fn public_router_exposes_wallet_data_but_not_management_api() { 1099 let (app, wallet, _, _dir) = test_app().await; 1100 let receive_address = encode_address(wallet.address(), AddressNetwork::Mainnet).unwrap(); 1101 let response = app 1102 .clone() 1103 .oneshot( 1104 Request::builder() 1105 .uri(format!("/v1/addresses/{receive_address}/balance")) 1106 .body(Body::empty()) 1107 .unwrap(), 1108 ) 1109 .await 1110 .unwrap(); 1111 assert_eq!(response.status(), StatusCode::OK); 1112 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap(); 1113 let value: serde_json::Value = serde_json::from_slice(&body).unwrap(); 1114 assert_eq!(value["confirmed"], 5 * MICRO_IUNA); 1115 assert_eq!(value["spendable"], 5 * MICRO_IUNA); 1116 assert_eq!(value["address"], receive_address); 1117 assert_eq!(value["public_key"], wallet.address()); 1118 1119 let response = app 1120 .clone() 1121 .oneshot( 1122 Request::builder() 1123 .method(Method::POST) 1124 .uri("/v1/wallets/snapshot") 1125 .header(header::CONTENT_TYPE, "application/json") 1126 .body(Body::from( 1127 serde_json::to_vec(&serde_json::json!({ 1128 "addresses": [receive_address.clone()] 1129 })) 1130 .unwrap(), 1131 )) 1132 .unwrap(), 1133 ) 1134 .await 1135 .unwrap(); 1136 assert_eq!(response.status(), StatusCode::OK); 1137 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap(); 1138 let value: serde_json::Value = serde_json::from_slice(&body).unwrap(); 1139 assert_eq!(value["confirmed"], 5 * MICRO_IUNA); 1140 assert_eq!(value["spendable"], 5 * MICRO_IUNA); 1141 assert_eq!(value["addresses"][0]["version"], 0); 1142 assert_eq!(value["addresses"][0]["used"], false); 1143 assert_eq!(value["utxos"].as_array().unwrap().len(), 1); 1144 1145 let wrong_network = encode_address(wallet.address(), AddressNetwork::Testnet).unwrap(); 1146 let response = app 1147 .clone() 1148 .oneshot( 1149 Request::builder() 1150 .uri(format!("/v1/addresses/{wrong_network}/balance")) 1151 .body(Body::empty()) 1152 .unwrap(), 1153 ) 1154 .await 1155 .unwrap(); 1156 assert_eq!(response.status(), StatusCode::BAD_REQUEST); 1157 1158 let response = app 1159 .oneshot( 1160 Request::builder() 1161 .uri("/api/config") 1162 .body(Body::empty()) 1163 .unwrap(), 1164 ) 1165 .await 1166 .unwrap(); 1167 assert_eq!(response.status(), StatusCode::NOT_FOUND); 1168 } 1169 1170 #[tokio::test] 1171 async fn signed_transfer_is_accepted_and_available_by_signature() { 1172 let (app, wallet, ledger, _dir) = test_app().await; 1173 let receive_address = encode_address(wallet.address(), AddressNetwork::Mainnet).unwrap(); 1174 let recipient = Wallet::from_seed("public-wallet-endpoint-recipient"); 1175 let transaction = ledger 1176 .build_transfer(&wallet, recipient.address(), MICRO_IUNA, 1) 1177 .unwrap(); 1178 let transaction_id = transaction.signature().to_string(); 1179 let response = app 1180 .clone() 1181 .oneshot( 1182 Request::builder() 1183 .method(Method::POST) 1184 .uri("/v1/transactions") 1185 .header(header::CONTENT_TYPE, "application/json") 1186 .body(Body::from(serde_json::to_vec(&transaction).unwrap())) 1187 .unwrap(), 1188 ) 1189 .await 1190 .unwrap(); 1191 assert_eq!(response.status(), StatusCode::ACCEPTED); 1192 1193 let response = app 1194 .clone() 1195 .oneshot( 1196 Request::builder() 1197 .uri(format!("/v1/transactions/{transaction_id}")) 1198 .body(Body::empty()) 1199 .unwrap(), 1200 ) 1201 .await 1202 .unwrap(); 1203 assert_eq!(response.status(), StatusCode::OK); 1204 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap(); 1205 let value: serde_json::Value = serde_json::from_slice(&body).unwrap(); 1206 assert_eq!(value["status"], "pending"); 1207 1208 let response = app 1209 .clone() 1210 .oneshot( 1211 Request::builder() 1212 .uri(format!("/v1/addresses/{receive_address}/transactions")) 1213 .body(Body::empty()) 1214 .unwrap(), 1215 ) 1216 .await 1217 .unwrap(); 1218 assert_eq!(response.status(), StatusCode::OK); 1219 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap(); 1220 let value: serde_json::Value = serde_json::from_slice(&body).unwrap(); 1221 assert_eq!(value["total"], 1); 1222 assert_eq!(value["items"][0]["status"], "pending"); 1223 1224 let response = app 1225 .clone() 1226 .oneshot( 1227 Request::builder() 1228 .method(Method::POST) 1229 .uri("/v1/wallets/transactions") 1230 .header(header::CONTENT_TYPE, "application/json") 1231 .body(Body::from( 1232 serde_json::to_vec(&serde_json::json!({ 1233 "addresses": [receive_address.clone()], 1234 "limit": 10 1235 })) 1236 .unwrap(), 1237 )) 1238 .unwrap(), 1239 ) 1240 .await 1241 .unwrap(); 1242 assert_eq!(response.status(), StatusCode::OK); 1243 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap(); 1244 let value: serde_json::Value = serde_json::from_slice(&body).unwrap(); 1245 assert_eq!(value["total"], 1); 1246 assert_eq!(value["items"][0]["status"], "pending"); 1247 assert_eq!(value["items"][0]["direction"], "sent"); 1248 1249 let response = app 1250 .oneshot( 1251 Request::builder() 1252 .uri(format!( 1253 "/v1/addresses/{receive_address}/transactions?tx=false&mine=true&burn=false&reward=false" 1254 )) 1255 .body(Body::empty()) 1256 .unwrap(), 1257 ) 1258 .await 1259 .unwrap(); 1260 assert_eq!(response.status(), StatusCode::OK); 1261 let body = to_bytes(response.into_body(), 64 * 1024).await.unwrap(); 1262 let value: serde_json::Value = serde_json::from_slice(&body).unwrap(); 1263 assert_eq!(value["total"], 0); 1264 assert_eq!(value["has_more"], false); 1265 assert_eq!(value["items"].as_array().unwrap().len(), 0); 1266 } 1267 }