diff --git a/Cargo.lock b/Cargo.lock index 53da60c..96ba2e6 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1704,7 +1704,6 @@ dependencies = [ "log", "mostro-core", "nostr-sdk", - "openssl", "pretty_env_logger", "reqwest", "serde", @@ -1943,15 +1942,6 @@ version = "0.1.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ff011a302c396a5197692431fc1948019154afc178baf7d8e37367442a4601cf" -[[package]] -name = "openssl-src" -version = "300.4.2+3.4.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "168ce4e058f975fe43e89d9ccf78ca668601887ae736090aacc23ae353c298e2" -dependencies = [ - "cc", -] - [[package]] name = "openssl-sys" version = "0.9.104" @@ -1960,7 +1950,6 @@ checksum = "45abf306cbf99debc8195b66b7346498d7b10c210de50418b5ccd7ceba08c741" dependencies = [ "cc", "libc", - "openssl-src", "pkg-config", "vcpkg", ] diff --git a/Cargo.toml b/Cargo.toml index b8bdc21..4bf41bb 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -42,7 +42,7 @@ reqwest = { version = "0.12.4", features = ["json"] } mostro-core = "0.6.40" lnurl-rs = "0.9.0" pretty_env_logger = "0.5.0" -openssl = { version = "0.10.68", features = ["vendored"] } +# openssl = { version = "0.10.68", features = ["vendored"] } sqlx = { version = "0.8.2", features = ["sqlite", "runtime-tokio-native-tls"] } bip39 = { version = "2.1.0", features = ["rand"] } dirs = "5.0.1" \ No newline at end of file diff --git a/src/cli.rs b/src/cli.rs index 933c463..25f6afd 100644 --- a/src/cli.rs +++ b/src/cli.rs @@ -19,7 +19,6 @@ use crate::cli::list_orders::execute_list_orders; use crate::cli::new_order::execute_new_order; use crate::cli::rate_user::execute_rate_user; use crate::cli::send_dm::execute_send_dm; -use crate::cli::send_msg::execute_send_msg; use crate::cli::take_buy::execute_take_buy; use crate::cli::take_dispute::execute_take_dispute; use crate::cli::take_sell::execute_take_sell; @@ -29,6 +28,8 @@ use crate::util; use anyhow::{Error, Result}; use clap::{Parser, Subcommand}; use nostr_sdk::prelude::*; +use sqlx::SqlitePool; +use std::sync::OnceLock; use std::{ env::{set_var, var}, str::FromStr, @@ -36,6 +37,22 @@ use std::{ use take_dispute::*; use uuid::Uuid; +pub static IDENTITY_KEYS: OnceLock = OnceLock::new(); +pub static MOSTRO_KEYS: OnceLock = OnceLock::new(); +pub static MOSTRO_PUBKEY: OnceLock = OnceLock::new(); +pub static POOL: OnceLock = OnceLock::new(); +pub static TRADE_KEY: OnceLock<(Keys, i64)> = OnceLock::new(); + +pub struct Context { + pub client: Client, + pub identity_keys: Keys, + pub trade_keys: Keys, + pub trade_index: i64, + pub pool: SqlitePool, + pub mostro_keys: Keys, + pub mostro_pubkey: PublicKey, +} + #[derive(Parser)] #[command( name = "mostro-cli", @@ -245,6 +262,31 @@ pub enum Commands { }, } +fn get_env_var(cli: &Cli) { + // Init logger + if cli.verbose { + set_var("RUST_LOG", "info"); + pretty_env_logger::init(); + } + + if cli.mostropubkey.is_some() { + set_var("MOSTRO_PUBKEY", cli.mostropubkey.clone().unwrap()); + } + let _pubkey = var("MOSTRO_PUBKEY").expect("$MOSTRO_PUBKEY env var needs to be set"); + + if cli.relays.is_some() { + set_var("RELAYS", cli.relays.clone().unwrap()); + } + + if cli.pow.is_some() { + set_var("POW", cli.pow.clone().unwrap()); + } + + if cli.secret { + set_var("SECRET", "true"); + } +} + // Check range with two values value fn check_fiat_range(s: &str) -> Result<(i64, Option)> { if s.contains('-') { @@ -288,54 +330,82 @@ fn check_fiat_range(s: &str) -> Result<(i64, Option)> { pub async fn run() -> Result<()> { let cli = Cli::parse(); - // Init logger - if cli.verbose { - set_var("RUST_LOG", "info"); - pretty_env_logger::init(); - } + let ctx = init_context(&cli).await?; - if cli.mostropubkey.is_some() { - set_var("MOSTRO_PUBKEY", cli.mostropubkey.unwrap()); + if let Some(cmd) = &cli.command { + cmd.run(&ctx).await?; } - let pubkey = var("MOSTRO_PUBKEY").expect("$MOSTRO_PUBKEY env var needs to be set"); - if cli.relays.is_some() { - set_var("RELAYS", cli.relays.unwrap()); - } + println!("Bye Bye!"); - if cli.pow.is_some() { - set_var("POW", cli.pow.unwrap()); - } + Ok(()) +} - if cli.secret { - set_var("SECRET", "true"); - } +async fn init_context(cli: &Cli) -> Result { + // Get environment variables + get_env_var(cli); + // Initialize database pool let pool = connect().await?; + POOL.get_or_init(|| pool.clone()); + + // Get identity keys let identity_keys = User::get_identity_keys(&pool) .await .map_err(|e| anyhow::anyhow!("Failed to get identity keys: {}", e))?; + IDENTITY_KEYS.get_or_init(|| identity_keys.clone()); + // Get trade keys let (trade_keys, trade_index) = User::get_next_trade_keys(&pool) .await .map_err(|e| anyhow::anyhow!("Failed to get trade keys: {}", e))?; + TRADE_KEY.get_or_init(|| (trade_keys.clone(), trade_index)); - // Mostro pubkey - let mostro_key = PublicKey::from_str(&pubkey)?; + // Get Mostro admin keys + let mostro_keys = Keys::from_str( + &std::env::var("NSEC_PRIVKEY") + .map_err(|e| anyhow::anyhow!("Failed to get mostro keys: {}", e))?, + )?; + MOSTRO_KEYS.get_or_init(|| mostro_keys.clone()); + MOSTRO_PUBKEY.get_or_init(|| mostro_keys.public_key()); - // Call function to connect to relays + // Connect to Nostr relays let client = util::connect_nostr().await?; - if let Some(cmd) = cli.command { - match &cmd { + Ok(Context { + client, + identity_keys, + trade_keys, + trade_index, + pool, + mostro_keys: mostro_keys.clone(), + mostro_pubkey: mostro_keys.public_key(), + }) +} + +impl Commands { + pub async fn run(&self, ctx: &Context) -> Result<()> { + match self { Commands::ConversationKey { pubkey } => { - execute_conversation_key(&trade_keys, PublicKey::from_str(pubkey)?).await? + execute_conversation_key(&ctx.trade_keys, PublicKey::from_str(pubkey)?).await } Commands::ListOrders { status, currency, kind, - } => execute_list_orders(kind, currency, status, mostro_key, &client).await?, + } => { + execute_list_orders( + kind, + currency, + status, + ctx.mostro_pubkey, + &ctx.mostro_keys, + ctx.trade_index, + &ctx.pool, + &ctx.client, + ) + .await + } Commands::TakeSell { order_id, invoice, @@ -345,58 +415,80 @@ pub async fn run() -> Result<()> { order_id, invoice, *amount, - &identity_keys, - &trade_keys, - trade_index, - mostro_key, - &client, + &ctx.identity_keys, + &ctx.trade_keys, + ctx.trade_index, + ctx.mostro_pubkey, + &ctx.client, ) - .await? + .await } Commands::TakeBuy { order_id, amount } => { execute_take_buy( order_id, *amount, - &identity_keys, - &trade_keys, - trade_index, - mostro_key, - &client, + &ctx.identity_keys, + &ctx.trade_keys, + ctx.trade_index, + ctx.mostro_pubkey, + &ctx.client, ) - .await? + .await } Commands::AddInvoice { order_id, invoice } => { - execute_add_invoice(order_id, invoice, &identity_keys, mostro_key, &client).await? + execute_add_invoice( + order_id, + invoice, + &ctx.identity_keys, + ctx.mostro_pubkey, + &ctx.client, + ) + .await } Commands::GetDm { since, from_user } => { - execute_get_dm(since, trade_index, &client, *from_user, false).await? + execute_get_dm( + since, + ctx.trade_index, + &ctx.mostro_keys, + &ctx.client, + *from_user, + false, + ) + .await } Commands::GetAdminDm { since, from_user } => { - execute_get_dm(since, trade_index, &client, *from_user, true).await? + execute_get_dm( + since, + ctx.trade_index, + &ctx.mostro_keys, + &ctx.client, + *from_user, + true, + ) + .await } Commands::FiatSent { order_id } | Commands::Release { order_id } | Commands::Dispute { order_id } | Commands::Cancel { order_id } => { - execute_send_msg( - cmd.clone(), - Some(*order_id), - Some(&identity_keys), - mostro_key, - &client, - None, + crate::util::run_simple_order_msg( + self.clone(), + order_id, + &ctx.identity_keys, + ctx.mostro_pubkey, + &ctx.client, ) - .await? + .await } Commands::AdmAddSolver { npubkey } => { - let id_key = match std::env::var("NSEC_PRIVKEY") { - Ok(id_key) => Keys::parse(&id_key)?, - Err(e) => { - println!("Failed to get mostro admin private key: {}", e); - std::process::exit(1); - } - }; - execute_admin_add_solver(npubkey, &id_key, &trade_keys, mostro_key, &client).await? + execute_admin_add_solver( + npubkey, + &ctx.mostro_keys, + &ctx.trade_keys, + ctx.mostro_pubkey, + &ctx.client, + ) + .await } Commands::NewOrder { kind, @@ -416,64 +508,73 @@ pub async fn run() -> Result<()> { payment_method, premium, invoice, - &identity_keys, - &trade_keys, - trade_index, - mostro_key, - &client, + &ctx.identity_keys, + &ctx.trade_keys, + ctx.trade_index, + ctx.mostro_pubkey, + &ctx.client, expiration_days, ) - .await? + .await } Commands::Rate { order_id, rating } => { - execute_rate_user(order_id, rating, &identity_keys, mostro_key, &client).await?; + execute_rate_user( + order_id, + rating, + &ctx.identity_keys, + ctx.mostro_pubkey, + &ctx.client, + ) + .await } Commands::AdmSettle { order_id } => { - let id_key = match std::env::var("NSEC_PRIVKEY") { - Ok(id_key) => Keys::parse(&id_key)?, - Err(e) => { - println!("Failed to get mostro admin private key: {}", e); - std::process::exit(1); - } - }; - execute_admin_settle_dispute(order_id, &id_key, &trade_keys, mostro_key, &client) - .await?; + execute_admin_settle_dispute( + order_id, + &ctx.mostro_keys, + &ctx.trade_keys, + ctx.mostro_pubkey, + &ctx.client, + ) + .await } Commands::AdmCancel { order_id } => { - let id_key = match std::env::var("NSEC_PRIVKEY") { - Ok(id_key) => Keys::parse(&id_key)?, - Err(e) => { - println!("Failed to get mostro admin private key: {}", e); - std::process::exit(1); - } - }; - execute_admin_cancel_dispute(order_id, &id_key, &trade_keys, mostro_key, &client) - .await?; + execute_admin_cancel_dispute( + order_id, + &ctx.mostro_keys, + &ctx.trade_keys, + ctx.mostro_pubkey, + &ctx.client, + ) + .await } Commands::AdmTakeDispute { dispute_id } => { - let id_key = match std::env::var("NSEC_PRIVKEY") { - Ok(id_key) => Keys::parse(&id_key)?, - Err(e) => { - println!("Failed to get mostro admin private key: {}", e); - std::process::exit(1); - } - }; - - execute_take_dispute(dispute_id, &id_key, &trade_keys, mostro_key, &client).await? + execute_take_dispute( + dispute_id, + &ctx.mostro_keys, + &ctx.trade_keys, + ctx.mostro_pubkey, + &ctx.client, + ) + .await + } + Commands::AdmListDisputes {} => { + execute_list_disputes( + ctx.mostro_pubkey, + &ctx.mostro_keys, + ctx.trade_index, + &ctx.pool, + &ctx.client, + ) + .await } - Commands::AdmListDisputes {} => execute_list_disputes(mostro_key, &client).await?, Commands::SendDm { pubkey, order_id, message, } => { let pubkey = PublicKey::from_str(pubkey)?; - execute_send_dm(pubkey, &client, order_id, message).await? + execute_send_dm(pubkey, &ctx.client, order_id, message).await } - }; + } } - - println!("Bye Bye!"); - - Ok(()) } diff --git a/src/cli/add_invoice.rs b/src/cli/add_invoice.rs index e9ee8ef..cbbe912 100644 --- a/src/cli/add_invoice.rs +++ b/src/cli/add_invoice.rs @@ -1,5 +1,5 @@ use crate::db::connect; -use crate::util::send_message_sync; +use crate::util::{send_dm, wait_for_dm}; use crate::{db::Order, lightning::is_valid_invoice}; use anyhow::Result; use lnurl::lightning_address::LightningAddress; @@ -16,7 +16,7 @@ pub async fn execute_add_invoice( client: &Client, ) -> Result<()> { let pool = connect().await?; - let mut order = Order::get_by_id(&pool, &order_id.to_string()).await?; + let order = Order::get_by_id(&pool, &order_id.to_string()).await?; let trade_keys = order .trade_keys .clone() @@ -50,31 +50,35 @@ pub async fn execute_add_invoice( payload, ); - let dm = send_message_sync( - client, - Some(identity_keys), - &trade_keys, - mostro_key, - add_invoice_message, - true, - false, - ) - .await?; + let message_json = add_invoice_message + .as_json() + .map_err(|_| anyhow::anyhow!("Failed to serialize message"))?; - dm.iter().for_each(|el| { - let message = el.0.get_inner_message_kind(); - if message.request_id == Some(request_id) && message.action == Action::WaitingSellerToPay { - println!("Now we should wait for the seller to pay the invoice"); + // Clone the keys and client for the async call + let identity_keys = identity_keys.clone(); + let trade_keys_clone = trade_keys.clone(); + let client_clone = client.clone(); + + // Spawn a new task to send the DM + // This is so we can wait for the gift wrap event in the main thread + tokio::spawn(async move { + if let Err(e) = send_dm( + &client_clone, + Some(&identity_keys.clone()), + &trade_keys_clone, + &mostro_key, + message_json, + None, + false, + ) + .await + { + eprintln!("Failed to send DM: {}", e); } }); - match order - .set_status(Status::WaitingPayment.to_string()) - .save(&pool) - .await - { - Ok(_) => println!("Order status updated"), - Err(e) => println!("Failed to update order status: {}", e), - } + + // Wait for the DM to be sent from mostro + wait_for_dm(client, &trade_keys, request_id, 0, Some(order)).await?; Ok(()) } diff --git a/src/cli/get_dm.rs b/src/cli/get_dm.rs index 1e05329..fa2d62e 100644 --- a/src/cli/get_dm.rs +++ b/src/cli/get_dm.rs @@ -5,14 +5,16 @@ use nostr_sdk::prelude::*; use crate::{ db::{connect, Order, User}, - util::get_direct_messages, + parser::dms::parse_dm_events, + util::{create_filter, ListKind}, }; pub async fn execute_get_dm( - since: &i64, + _since: &i64, trade_index: i64, + mostro_keys: &Keys, client: &Client, - from_user: bool, + _from_user: bool, admin: bool, ) -> Result<()> { let mut dm: Vec<(Message, u64)> = Vec::new(); @@ -20,18 +22,19 @@ pub async fn execute_get_dm( if !admin { for index in 1..=trade_index { let keys = User::get_trade_keys(&pool, index).await?; - let dm_temp = get_direct_messages(client, &keys, *since, from_user).await; + let filter = create_filter(ListKind::DirectMessagesUser, keys.public_key()); + let fetched_events = client + .fetch_events(filter, std::time::Duration::from_secs(15)) + .await?; + let dm_temp = parse_dm_events(fetched_events, &keys).await; dm.extend(dm_temp); } } else { - let id_key = match std::env::var("NSEC_PRIVKEY") { - Ok(id_key) => Keys::parse(&id_key)?, - Err(e) => { - println!("Failed to get mostro admin private key: {}", e); - std::process::exit(1); - } - }; - let dm_temp = get_direct_messages(client, &id_key, *since, from_user).await; + let filter = create_filter(ListKind::DirectMessagesAdmin, mostro_keys.public_key()); + let fetched_events = client + .fetch_events(filter, std::time::Duration::from_secs(15)) + .await?; + let dm_temp = parse_dm_events(fetched_events, mostro_keys).await; dm.extend(dm_temp); } diff --git a/src/cli/list_disputes.rs b/src/cli/list_disputes.rs index 6951647..47c7056 100644 --- a/src/cli/list_disputes.rs +++ b/src/cli/list_disputes.rs @@ -1,17 +1,35 @@ use anyhow::Result; use nostr_sdk::prelude::*; +use sqlx::SqlitePool; -use crate::pretty_table::print_disputes_table; -use crate::util::get_disputes_list; +use crate::parser::disputes::print_disputes_table; +use crate::util::{fetch_events_list, ListKind}; -pub async fn execute_list_disputes(mostro_key: PublicKey, client: &Client) -> Result<()> { +pub async fn execute_list_disputes( + mostro_key: PublicKey, + mostro_keys: &Keys, + trade_index: i64, + pool: &SqlitePool, + client: &Client, +) -> Result<()> { println!( "Requesting disputes from mostro pubId - {}", mostro_key.clone() ); // Get orders from relays - let table_of_disputes = get_disputes_list(mostro_key, client).await?; + let table_of_disputes = fetch_events_list( + ListKind::Disputes, + None, + None, + None, + mostro_key, + mostro_keys, + trade_index, + pool, + client, + ) + .await?; let table = print_disputes_table(table_of_disputes)?; println!("{table}"); diff --git a/src/cli/list_orders.rs b/src/cli/list_orders.rs index ee77bb6..ac2cccd 100644 --- a/src/cli/list_orders.rs +++ b/src/cli/list_orders.rs @@ -1,16 +1,20 @@ +use crate::parser::orders::print_orders_table; +use crate::util::{fetch_events_list, ListKind}; use anyhow::Result; use mostro_core::prelude::*; use nostr_sdk::prelude::*; +use sqlx::SqlitePool; use std::str::FromStr; -use crate::pretty_table::print_orders_table; -use crate::util::get_orders_list; - +#[allow(clippy::too_many_arguments)] pub async fn execute_list_orders( kind: &Option, currency: &Option, status: &Option, - mostro_key: PublicKey, + mostro_pubkey: PublicKey, + mostro_keys: &Keys, + trade_index: i64, + pool: &SqlitePool, client: &Client, ) -> Result<()> { // Used to get upper currency string to check against a list of tickers @@ -29,7 +33,9 @@ pub async fn execute_list_orders( ); // New check against strings if let Some(k) = kind { - kind_checked = Some(mostro_core::order::Kind::from_str(k).expect("Not valid order kind! Please check")); + kind_checked = Some( + mostro_core::order::Kind::from_str(k).expect("Not valid order kind! Please check"), + ); println!("You are searching {} orders", kind_checked.unwrap()); } @@ -42,17 +48,18 @@ pub async fn execute_list_orders( ); } - println!( - "Requesting orders from mostro pubId - {}", - mostro_key.clone() - ); + println!("Requesting orders from mostro pubId - {}", mostro_pubkey); // Get orders from relays - let table_of_orders = get_orders_list( - mostro_key, - status_checked.unwrap(), + let table_of_orders = fetch_events_list( + ListKind::Orders, + status_checked, upper_currency, kind_checked, + mostro_pubkey, + mostro_keys, + trade_index, + pool, client, ) .await?; diff --git a/src/cli/new_order.rs b/src/cli/new_order.rs index be1666d..67b1c41 100644 --- a/src/cli/new_order.rs +++ b/src/cli/new_order.rs @@ -7,9 +7,8 @@ use std::process; use std::str::FromStr; use uuid::Uuid; -use crate::db::{connect, Order, User}; -use crate::pretty_table::print_order_preview; -use crate::util::{send_message_sync, uppercase_first}; +use crate::parser::orders::print_order_preview; +use crate::util::{send_dm, uppercase_first, wait_for_dm}; pub type FiatNames = HashMap; @@ -121,77 +120,50 @@ pub async fn execute_new_order( Some(order_content), ); - let dm = send_message_sync( - client, - Some(identity_keys), - trade_keys, - mostro_key, - message, - true, - false, - ) - .await?; - let order_id = dm - .iter() - .find_map(|el| { - let message = el.0.get_inner_message_kind(); - if message.request_id == Some(request_id) { - match message.action { - Action::NewOrder => { - if let Some(Payload::Order(order)) = message.payload.as_ref() { - return order.id; - } - } - Action::CantDo => { - if let Some(Payload::CantDo(Some(cant_do_reason))) = &message.payload { - match cant_do_reason { - CantDoReason::OutOfRangeFiatAmount | CantDoReason::OutOfRangeSatsAmount => { - println!("Error: Amount is outside the allowed range. Please check the order's min/max limits."); - } - _ => { - println!("Unknown reason: {:?}", message.payload); - } - } - } else { - println!("Unknown reason: {:?}", message.payload); - return None; - } - } - _ => { - println!("Unknown action: {:?}", message.action); - return None; - } - } - } - None - }) - .or_else(|| { - println!("Error: No matching order found in response"); - None - }); - - if let Some(order_id) = order_id { - println!("Order id {} created", order_id); - // Create order in db - let pool = connect().await?; - let db_order = Order::new(&pool, small_order, trade_keys, Some(request_id as i64)) - .await - .map_err(|e| anyhow::anyhow!("Failed to create DB order: {:?}", e))?; - // Update last trade index - match User::get(&pool).await { - Ok(mut user) => { - user.set_last_trade_index(trade_index); - if let Err(e) = user.save(&pool).await { - println!("Failed to update user: {}", e); - } - } - Err(e) => println!("Failed to get user: {}", e), - } - let db_order_id = db_order - .id - .clone() - .ok_or(anyhow::anyhow!("Missing order id"))?; - Order::save_new_id(&pool, db_order_id, order_id.to_string()).await?; - } + // Send dm to receiver pubkey + println!( + "SENDING DM with trade keys: {:?}", + trade_keys.public_key().to_hex() + ); + + // Serialize the message + let message_json = message + .as_json() + .map_err(|_| anyhow::anyhow!("Failed to serialize message"))?; + + // Clone the keys and client for the async call + let identity_keys = identity_keys.clone(); + let trade_keys_clone = trade_keys.clone(); + // let mostro_key = mostro_key.clone(); + let client_clone = client.clone(); + + // Subscribe to gift wrap events - ONLY NEW ONES WITH LIMIT 0 + let subscription = Filter::new() + .pubkey(trade_keys.public_key()) + .kind(nostr_sdk::Kind::GiftWrap) + .limit(0); + + let opts = SubscribeAutoCloseOptions::default().exit_policy(ReqExitPolicy::WaitForEvents(1)); + + client.subscribe(subscription, Some(opts)).await?; + + // Spawn a new task to send the DM + // This is so we can wait for the gift wrap event in the main thread + tokio::spawn(async move { + let _ = send_dm( + &client_clone, + Some(&identity_keys.clone()), + &trade_keys_clone, + &mostro_key, + message_json, + None, + false, + ) + .await; + }); + + // Wait for the DM to be sent from mostro + wait_for_dm(client, trade_keys, request_id, trade_index, None).await?; + Ok(()) } diff --git a/src/cli/send_msg.rs b/src/cli/send_msg.rs index be6f2e3..8eaf05f 100644 --- a/src/cli/send_msg.rs +++ b/src/cli/send_msg.rs @@ -1,5 +1,5 @@ use crate::db::{Order, User}; -use crate::util::send_message_sync; +use crate::util::wait_for_dm; use crate::{cli::Commands, db::connect}; use anyhow::Result; @@ -34,7 +34,9 @@ pub async fn execute_send_msg( println!( "Sending {} command for order {:?} to mostro pubId {}", - requested_action, order_id, mostro_key + requested_action, + order_id.as_ref(), + mostro_key ); let pool = connect().await?; @@ -62,21 +64,55 @@ pub async fn execute_send_msg( // Create and send the message let message = Message::new_order(order_id, Some(request_id), None, requested_action, payload); - // println!("Sending message: {:#?}", message); + let client_clone = client.clone(); + let idkey = identity_keys + .ok_or_else(|| anyhow::anyhow!("Identity keys are required"))? + .to_owned(); if let Some(order_id) = order_id { - handle_order_response( - &pool, - client, - identity_keys, - mostro_key, - message, - order_id, - request_id, - ) - .await?; - } else { - println!("Error: Missing order ID"); + let order = Order::get_by_id(&pool, &order_id.to_string()).await?; + + if let Some(trade_keys_str) = order.trade_keys.clone() { + let trade_keys = Keys::parse(&trade_keys_str)?; + // Subscribe to gift wrap events - ONLY NEW ONES WITH LIMIT 0 + let subscription = Filter::new() + .pubkey(trade_keys.public_key()) + .kind(nostr_sdk::Kind::GiftWrap) + .limit(0); + + let opts = + SubscribeAutoCloseOptions::default().exit_policy(ReqExitPolicy::WaitForEvents(1)); + + client.subscribe(subscription, Some(opts)).await?; + // Clone the keys and client for the async call + let trade_keys_clone = trade_keys.clone(); + + // Spawn a new task to send the DM + // This is so we can wait for the gift wrap event in the main thread + tokio::spawn(async move { + match message.as_json() { + Ok(message_json) => { + if let Err(e) = crate::util::send_dm( + &client_clone, + Some(&idkey), + &trade_keys_clone, + &mostro_key, + message_json, + None, + false, + ) + .await + { + eprintln!("Failed to send DM: {}", e); + } + } + Err(e) => eprintln!("Failed to serialize message: {}", e), + } + }); + + // Wait for the DM to be sent from mostro + wait_for_dm(client, &trade_keys, request_id, 0, Some(order)).await?; + } } Ok(()) @@ -103,84 +139,3 @@ async fn create_next_trade_payload( } Ok(None) } - -async fn handle_order_response( - pool: &SqlitePool, - client: &Client, - identity_keys: Option<&Keys>, - mostro_key: PublicKey, - message: Message, - order_id: Uuid, - request_id: u64, -) -> Result<()> { - let order = Order::get_by_id(pool, &order_id.to_string()).await; - - match order { - Ok(order) => { - if let Some(trade_keys_str) = order.trade_keys { - let trade_keys = Keys::parse(&trade_keys_str)?; - let dm = send_message_sync( - client, - identity_keys, - &trade_keys, - mostro_key, - message, - true, - false, - ) - .await?; - process_order_response(dm, pool, &trade_keys, request_id).await?; - } else { - println!("Error: Missing trade keys for order {}", order_id); - } - } - Err(e) => { - println!("Error: {}", e); - } - } - - Ok(()) -} - -async fn process_order_response( - dm: Vec<(Message, u64)>, - pool: &SqlitePool, - trade_keys: &Keys, - request_id: u64, -) -> Result<()> { - for (message, _) in dm { - let kind = message.get_inner_message_kind(); - if let Some(req_id) = kind.request_id { - if req_id != request_id { - continue; - } - - match kind.action { - Action::NewOrder => { - if let Some(Payload::Order(order)) = kind.payload.as_ref() { - Order::new(pool, order.clone(), trade_keys, Some(request_id as i64)) - .await - .map_err(|e| anyhow::anyhow!("Failed to create new order: {}", e))?; - return Ok(()); - } - } - Action::Canceled => { - if let Some(id) = kind.id { - // Verify order exists before deletion - if Order::get_by_id(pool, &id.to_string()).await.is_ok() { - Order::delete_by_id(pool, &id.to_string()) - .await - .map_err(|e| anyhow::anyhow!("Failed to delete order: {}", e))?; - return Ok(()); - } else { - return Err(anyhow::anyhow!("Order not found: {}", id)); - } - } - } - _ => (), - } - } - } - - Ok(()) -} diff --git a/src/cli/take_buy.rs b/src/cli/take_buy.rs index 55abc8a..883003c 100644 --- a/src/cli/take_buy.rs +++ b/src/cli/take_buy.rs @@ -3,10 +3,7 @@ use mostro_core::prelude::*; use nostr_sdk::prelude::*; use uuid::Uuid; -use crate::{ - db::{connect, Order, User}, - util::send_message_sync, -}; +use crate::util::{send_dm, wait_for_dm}; pub async fn execute_take_buy( order_id: &Uuid, @@ -33,79 +30,47 @@ pub async fn execute_take_buy( payload, ); - let dm = send_message_sync( - client, - Some(identity_keys), - trade_keys, - mostro_key, - take_buy_message, - true, - false, - ) - .await?; + // Send dm to receiver pubkey + println!( + "SENDING DM with trade keys: {:?}", + trade_keys.public_key().to_hex() + ); - let pool = connect().await?; + let message_json = take_buy_message + .as_json() + .map_err(|_| anyhow::anyhow!("Failed to serialize message"))?; - let order = dm.iter().find_map(|el| { - let message = el.0.get_inner_message_kind(); - if message.request_id == Some(request_id) { - match message.action { - Action::PayInvoice => { - if let Some(Payload::PaymentRequest(order, invoice, _)) = &message.payload { - println!( - "Mostro sent you this hold invoice for order id: {}", - order - .as_ref() - .and_then(|o| o.id) - .map_or("unknown".to_string(), |id| id.to_string()) - ); - println!(); - println!("Pay this invoice to continue --> {}", invoice); - println!(); - return order.clone(); - } - } - Action::CantDo => { - if let Some(Payload::CantDo(Some(cant_do_reason))) = &message.payload { - match cant_do_reason { - CantDoReason::OutOfRangeFiatAmount | CantDoReason::OutOfRangeSatsAmount => { - println!("Error: Amount is outside the allowed range. Please check the order's min/max limits."); - } - _ => { - println!("Unknown reason: {:?}", message.payload); - } - } - } else { - println!("Unknown reason: {:?}", message.payload); - return None; - } - } - _ => { - println!("Unknown action: {:?}", message.action); - return None; - } - } - } - None + // Clone the keys and client for the async call + let identity_keys = identity_keys.clone(); + let trade_keys_clone = trade_keys.clone(); + let client_clone = client.clone(); + // Subscribe to gift wrap events - ONLY NEW ONES WITH LIMIT 0 + let subscription = Filter::new() + .pubkey(trade_keys.public_key()) + .kind(nostr_sdk::Kind::GiftWrap) + .limit(0); + // Subscribe to gift wrap events -waiting for 1 event + let opts = SubscribeAutoCloseOptions::default().exit_policy(ReqExitPolicy::WaitForEvents(1)); + //Activate the subscription + client.subscribe(subscription, Some(opts)).await?; + + // Spawn a new task to send the DM + // This is so we can wait for the gift wrap event in the main thread + tokio::spawn(async move { + let _ = send_dm( + &client_clone, + Some(&identity_keys.clone()), + &trade_keys_clone, + &mostro_key, + message_json, + None, + false, + ) + .await; }); - if let Some(o) = order { - match Order::new(&pool, o, trade_keys, Some(request_id as i64)).await { - Ok(order) => { - println!("Order {} created", order.id.unwrap()); - // Update last trade index to be used in next trade - match User::get(&pool).await { - Ok(mut user) => { - user.set_last_trade_index(trade_index); - if let Err(e) = user.save(&pool).await { - println!("Failed to update user: {}", e); - } - } - Err(e) => println!("Failed to get user: {}", e), - } - } - Err(e) => println!("{}", e), - } - } + + // Wait for the DM to be sent from mostro + wait_for_dm(client, trade_keys, request_id, trade_index, None).await?; Ok(()) } diff --git a/src/cli/take_sell.rs b/src/cli/take_sell.rs index b1a50b7..52b9bc4 100644 --- a/src/cli/take_sell.rs +++ b/src/cli/take_sell.rs @@ -6,9 +6,8 @@ use nostr_sdk::prelude::*; use std::str::FromStr; use uuid::Uuid; -use crate::db::{connect, Order, User}; use crate::lightning::is_valid_invoice; -use crate::util::send_message_sync; +use crate::util::{send_dm, wait_for_dm}; #[allow(clippy::too_many_arguments)] pub async fn execute_take_sell( @@ -65,73 +64,55 @@ pub async fn execute_take_sell( Some(payload), ); - let dm = send_message_sync( - client, - Some(identity_keys), - trade_keys, - mostro_key, - take_sell_message, - true, - false, - ) - .await?; - let pool = connect().await?; + // Send dm to receiver pubkey + println!( + "SENDING DM with trade keys: {:?}", + trade_keys.public_key().to_hex() + ); + let message_json = take_sell_message + .as_json() + .map_err(|_| anyhow::anyhow!("Failed to serialize message"))?; - let order = dm.iter().find_map(|el| { - let message = el.0.get_inner_message_kind(); - if message.request_id == Some(request_id) { - match message.action { - Action::AddInvoice => { - if let Some(Payload::Order(order)) = message.payload.as_ref() { - println!( - "Please add a lightning invoice with amount of {}", - order.amount - ); - return Some(order.clone()); - } - } - Action::CantDo => { - if let Some(Payload::CantDo(Some(cant_do_reason))) = &message.payload { - match cant_do_reason { - CantDoReason::OutOfRangeFiatAmount | CantDoReason::OutOfRangeSatsAmount => { - println!("Error: Amount is outside the allowed range. Please check the order's min/max limits."); - } - _ => { - println!("Unknown reason: {:?}", message.payload); - } - } - } else { - println!("Unknown reason: {:?}", message.payload); - return None; - } - } - _ => { - println!("Unknown action: {:?}", message.action); - return None; - } - } - } - None + // Clone the keys and client for the async call + let identity_keys = identity_keys.clone(); + let trade_keys_clone = trade_keys.clone(); + let client_clone = client.clone(); + // Subscribe to gift wrap events - ONLY NEW ONES WITH LIMIT 0 + let subscription = Filter::new() + .pubkey(trade_keys.public_key()) + .kind(nostr_sdk::Kind::GiftWrap) + .limit(0); + + let opts = SubscribeAutoCloseOptions::default().exit_policy(ReqExitPolicy::WaitForEvents(1)); + + client.subscribe(subscription, Some(opts)).await?; + + // Spawn a new task to send the DM + // This is so we can wait for the gift wrap event in the main thread + tokio::spawn(async move { + let _ = send_dm( + &client_clone, + Some(&identity_keys.clone()), + &trade_keys_clone, + &mostro_key, + message_json, + None, + false, + ) + .await; }); - if let Some(o) = order { - if let Ok(order) = Order::new(&pool, o, trade_keys, Some(request_id as i64)).await { - if let Some(order_id) = order.id { - println!("Order {} created", order_id); - } else { - println!("Warning: The newly created order has no ID."); - } - // Update last trade index to be used in next trade - match User::get(&pool).await { - Ok(mut user) => { - user.set_last_trade_index(trade_index); - if let Err(e) = user.save(&pool).await { - println!("Failed to update user: {}", e); - } - } - Err(e) => println!("Failed to get user: {}", e), - } - } - } + + // Subscribe to gift wrap events - ONLY NEW ONES WITH LIMIT 0 + let subscription = Filter::new() + .pubkey(trade_keys.public_key()) + .kind(nostr_sdk::Kind::GiftWrap) + .since(Timestamp::from(chrono::Utc::now().timestamp() as u64)) + .limit(2); + + client.subscribe(subscription, None).await?; + + // Wait for the DM to be sent from mostro + wait_for_dm(client, trade_keys, request_id, trade_index, None).await?; Ok(()) } diff --git a/src/db.rs b/src/db.rs index 5b7caaa..821ecdd 100644 --- a/src/db.rs +++ b/src/db.rs @@ -257,11 +257,11 @@ impl Order { sqlx::query( r#" - INSERT INTO orders (id, kind, status, amount, min_amount, max_amount, - fiat_code, fiat_amount, payment_method, premium, trade_keys, - counterparty_pubkey, is_mine, buyer_invoice, buyer_token, seller_token, - request_id, created_at, expires_at) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + INSERT INTO orders (id, kind, status, amount, min_amount, max_amount, + fiat_code, fiat_amount, payment_method, premium, trade_keys, + counterparty_pubkey, is_mine, buyer_invoice, buyer_token, seller_token, + request_id, created_at, expires_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) "#, ) .bind(&order.id) diff --git a/src/lib.rs b/src/lib.rs index ce39579..14986a5 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -3,5 +3,5 @@ pub mod db; pub mod error; pub mod lightning; pub mod nip33; -pub mod pretty_table; +pub mod parser; pub mod util; diff --git a/src/parser/disputes.rs b/src/parser/disputes.rs new file mode 100644 index 0000000..db2a3f0 --- /dev/null +++ b/src/parser/disputes.rs @@ -0,0 +1,116 @@ +use anyhow::Result; +use chrono::DateTime; +use comfy_table::presets::UTF8_FULL; +use comfy_table::*; +use log::info; +use mostro_core::prelude::*; +use nostr_sdk::prelude::*; + +use crate::util::Event; + +use crate::nip33::dispute_from_tags; + +pub fn parse_dispute_events(events: Events) -> Vec { + // Extracted Disputes List + let mut disputes_list = Vec::::new(); + + // Scan events to extract all disputes + for event in events.into_iter() { + if let Ok(mut dispute) = dispute_from_tags(event.tags) { + info!("Found Dispute id : {:?}", dispute.id); + // Get created at field from Nostr event + dispute.created_at = event.created_at.as_u64() as i64; + disputes_list.push(dispute.clone()); + } + } + + let buffer_dispute_list = disputes_list.clone(); + // Order all element ( orders ) received to filter - discard disaligned messages + // if an order has an older message with the state we received is discarded for the latest one + disputes_list.retain(|keep| { + !buffer_dispute_list + .iter() + .any(|x| x.id == keep.id && x.created_at > keep.created_at) + }); + + // Sort by id to remove duplicates + disputes_list.sort_by(|a, b| b.id.cmp(&a.id)); + disputes_list.dedup_by(|a, b| a.id == b.id); + + // Finally sort list by creation time + disputes_list.sort_by(|a, b| b.created_at.cmp(&a.created_at)); + disputes_list +} + +pub fn print_disputes_table(disputes_table: Vec) -> Result { + // Convert Event to Dispute + let disputes_table: Vec = disputes_table + .into_iter() + .filter_map(|event| { + if let Event::Dispute(dispute) = event { + Some(dispute) + } else { + None + } + }) + .collect(); + + // Create table + let mut table = Table::new(); + //Table rows + let mut rows: Vec = Vec::new(); + + if disputes_table.is_empty() { + table + .load_preset(UTF8_FULL) + .set_content_arrangement(ContentArrangement::Dynamic) + .set_width(160) + .set_header(vec![Cell::new("Sorry...") + .add_attribute(Attribute::Bold) + .set_alignment(CellAlignment::Center)]); + + // Single row for error + let mut r = Row::new(); + + r.add_cell( + Cell::new("No disputes found with requested parameters...") + .fg(Color::Red) + .set_alignment(CellAlignment::Center), + ); + + //Push single error row + rows.push(r); + } else { + table + .load_preset(UTF8_FULL) + .set_content_arrangement(ContentArrangement::Dynamic) + .set_width(160) + .set_header(vec![ + Cell::new("Dispute Id") + .add_attribute(Attribute::Bold) + .set_alignment(CellAlignment::Center), + Cell::new("Status") + .add_attribute(Attribute::Bold) + .set_alignment(CellAlignment::Center), + Cell::new("Created") + .add_attribute(Attribute::Bold) + .set_alignment(CellAlignment::Center), + ]); + + //Iterate to create table of orders + for single_dispute in disputes_table.into_iter() { + let date = DateTime::from_timestamp(single_dispute.created_at, 0); + + let r = Row::from(vec![ + Cell::new(single_dispute.id).set_alignment(CellAlignment::Center), + Cell::new(single_dispute.status.to_string()).set_alignment(CellAlignment::Center), + Cell::new(date.unwrap()), + ]); + rows.push(r); + } + } + + table.add_rows(rows); + + Ok(table.to_string()) +} diff --git a/src/parser/dms.rs b/src/parser/dms.rs new file mode 100644 index 0000000..a896f32 --- /dev/null +++ b/src/parser/dms.rs @@ -0,0 +1,77 @@ +use base64::engine::general_purpose; +use base64::Engine; +use mostro_core::prelude::*; +use nip44::v2::{decrypt_to_bytes, ConversationKey}; +use nostr_sdk::prelude::*; + +pub async fn parse_dm_events(events: Events, pubkey: &Keys) -> Vec<(Message, u64)> { + let mut id_list = Vec::::new(); + let mut direct_messages: Vec<(Message, u64)> = Vec::new(); + + for dm in events.iter() { + if !id_list.contains(&dm.id) { + id_list.push(dm.id); + + let (created_at, message) = match dm.kind { + nostr_sdk::Kind::GiftWrap => { + let unwrapped_gift = match nip59::extract_rumor(pubkey, dm).await { + Ok(u) => u, + Err(_) => { + println!("Error unwrapping gift"); + continue; + } + }; + let (message, _): (Message, Option) = + serde_json::from_str(&unwrapped_gift.rumor.content).unwrap(); + (unwrapped_gift.rumor.created_at, message) + } + nostr_sdk::Kind::PrivateDirectMessage => { + let ck = + if let Ok(ck) = ConversationKey::derive(pubkey.secret_key(), &dm.pubkey) { + ck + } else { + continue; + }; + let b64decoded_content = + match general_purpose::STANDARD.decode(dm.content.as_bytes()) { + Ok(b64decoded_content) => b64decoded_content, + Err(_) => { + continue; + } + }; + let unencrypted_content = match decrypt_to_bytes(&ck, &b64decoded_content) { + Ok(bytes) => bytes, + Err(_) => { + continue; + } + }; + let message_str = match String::from_utf8(unencrypted_content) { + Ok(s) => s, + Err(_) => { + continue; + } + }; + let message = match Message::from_json(&message_str) { + Ok(m) => m, + Err(_) => { + continue; + } + }; + (dm.created_at, message) + } + _ => continue, + }; + + let since_time = chrono::Utc::now() + .checked_sub_signed(chrono::Duration::minutes(30)) + .unwrap() + .timestamp() as u64; + if created_at.as_u64() < since_time { + continue; + } + direct_messages.push((message, created_at.as_u64())); + } + } + direct_messages.sort_by(|a, b| a.1.cmp(&b.1)); + direct_messages +} diff --git a/src/parser/mod.rs b/src/parser/mod.rs new file mode 100644 index 0000000..3bc1427 --- /dev/null +++ b/src/parser/mod.rs @@ -0,0 +1,7 @@ +pub mod disputes; +pub mod dms; +pub mod orders; + +pub use disputes::parse_dispute_events; +pub use dms::parse_dm_events; +pub use orders::parse_orders_events; diff --git a/src/pretty_table.rs b/src/parser/orders.rs similarity index 69% rename from src/pretty_table.rs rename to src/parser/orders.rs index 97be574..ee4e9d1 100644 --- a/src/pretty_table.rs +++ b/src/parser/orders.rs @@ -1,9 +1,82 @@ +use crate::util::Event; use anyhow::Result; use chrono::DateTime; use comfy_table::presets::UTF8_FULL; use comfy_table::*; +use log::{error, info}; use mostro_core::prelude::*; +use nostr_sdk::prelude::*; +use crate::nip33::order_from_tags; + +pub fn parse_orders_events( + events: Events, + currency: Option, + status: Option, + kind: Option, +) -> Vec { + // Extracted Orders List + let mut complete_events_list = Vec::::new(); + let mut requested_orders_list = Vec::::new(); + + // Scan events to extract all orders + for event in events.iter() { + let order = order_from_tags(event.tags.clone()); + + if order.is_err() { + error!("{order:?}"); + continue; + } + if let Ok(mut order) = order { + info!("Found Order id : {:?}", order.id.unwrap()); + + if order.id.is_none() { + info!("Order ID is none"); + continue; + } + + if order.kind.is_none() { + info!("Order kind is none"); + continue; + } + + if order.status.is_none() { + info!("Order status is none"); + continue; + } + + // Get created at field from Nostr event + order.created_at = Some(event.created_at.as_u64() as i64); + complete_events_list.push(order.clone()); + if order.status.ne(&status) { + continue; + } + + if currency.is_some() && order.fiat_code.ne(¤cy.clone().unwrap()) { + continue; + } + + if kind.is_some() && order.kind.ne(&kind) { + continue; + } + // Add just requested orders requested by filtering + requested_orders_list.push(order); + } + // Order all element ( orders ) received to filter - discard disaligned messages + // if an order has an older message with the state we received is discarded for the latest one + requested_orders_list.retain(|keep| { + !complete_events_list + .iter() + .any(|x| x.id == keep.id && x.created_at > keep.created_at) + }); + // Sort by id to remove duplicates + requested_orders_list.sort_by(|a, b| b.id.cmp(&a.id)); + requested_orders_list.dedup_by(|a, b| a.id == b.id); + } + // Finally sort list by creation time + requested_orders_list.sort_by(|a, b| b.created_at.cmp(&a.created_at)); + requested_orders_list +} pub fn print_order_preview(ord: Payload) -> Result { let single_order = match ord { @@ -42,10 +115,10 @@ pub fn print_order_preview(ord: Payload) -> Result { let r = Row::from(vec![ if let Some(k) = single_order.kind { match k { - Kind::Buy => Cell::new(k.to_string()) + mostro_core::order::Kind::Buy => Cell::new(k.to_string()) .fg(Color::Green) .set_alignment(CellAlignment::Center), - Kind::Sell => Cell::new(k.to_string()) + mostro_core::order::Kind::Sell => Cell::new(k.to_string()) .fg(Color::Red) .set_alignment(CellAlignment::Center), } @@ -78,8 +151,19 @@ pub fn print_order_preview(ord: Payload) -> Result { Ok(table.to_string()) } -pub fn print_orders_table(orders_table: Vec) -> Result { +pub fn print_orders_table(orders_table: Vec) -> Result { let mut table = Table::new(); + // Convert Event to SmallOrder + let orders_table: Vec = orders_table + .into_iter() + .filter_map(|event| { + if let Event::SmallOrder(order) = event { + Some(order) + } else { + None + } + }) + .collect(); //Table rows let mut rows: Vec = Vec::new(); @@ -143,10 +227,10 @@ pub fn print_orders_table(orders_table: Vec) -> Result { let r = Row::from(vec![ if let Some(k) = single_order.kind { match k { - Kind::Buy => Cell::new(k.to_string()) + mostro_core::order::Kind::Buy => Cell::new(k.to_string()) .fg(Color::Green) .set_alignment(CellAlignment::Center), - Kind::Sell => Cell::new(k.to_string()) + mostro_core::order::Kind::Sell => Cell::new(k.to_string()) .fg(Color::Red) .set_alignment(CellAlignment::Center), } @@ -186,64 +270,3 @@ pub fn print_orders_table(orders_table: Vec) -> Result { Ok(table.to_string()) } - -pub fn print_disputes_table(disputes_table: Vec) -> Result { - let mut table = Table::new(); - - //Table rows - let mut rows: Vec = Vec::new(); - - if disputes_table.is_empty() { - table - .load_preset(UTF8_FULL) - .set_content_arrangement(ContentArrangement::Dynamic) - .set_width(160) - .set_header(vec![Cell::new("Sorry...") - .add_attribute(Attribute::Bold) - .set_alignment(CellAlignment::Center)]); - - // Single row for error - let mut r = Row::new(); - - r.add_cell( - Cell::new("No disputes found with requested parameters...") - .fg(Color::Red) - .set_alignment(CellAlignment::Center), - ); - - //Push single error row - rows.push(r); - } else { - table - .load_preset(UTF8_FULL) - .set_content_arrangement(ContentArrangement::Dynamic) - .set_width(160) - .set_header(vec![ - Cell::new("Dispute Id") - .add_attribute(Attribute::Bold) - .set_alignment(CellAlignment::Center), - Cell::new("Status") - .add_attribute(Attribute::Bold) - .set_alignment(CellAlignment::Center), - Cell::new("Created") - .add_attribute(Attribute::Bold) - .set_alignment(CellAlignment::Center), - ]); - - //Iterate to create table of orders - for single_dispute in disputes_table.into_iter() { - let date = DateTime::from_timestamp(single_dispute.created_at, 0); - - let r = Row::from(vec![ - Cell::new(single_dispute.id).set_alignment(CellAlignment::Center), - Cell::new(single_dispute.status.to_string()).set_alignment(CellAlignment::Center), - Cell::new(date.unwrap()), - ]); - rows.push(r); - } - } - - table.add_rows(rows); - - Ok(table.to_string()) -} diff --git a/src/util.rs b/src/util.rs index 47d6568..a9d3c41 100644 --- a/src/util.rs +++ b/src/util.rs @@ -1,16 +1,272 @@ -use crate::nip33::{dispute_from_tags, order_from_tags}; - +use crate::cli::send_msg::execute_send_msg; +use crate::cli::Commands; +use crate::db::{connect, Order, User}; +use crate::parser::{parse_dispute_events, parse_dm_events, parse_orders_events}; use anyhow::{Error, Result}; use base64::engine::general_purpose; use base64::Engine; use dotenvy::var; -use log::{error, info}; use mostro_core::prelude::*; -use nip44::v2::{decrypt_to_bytes, encrypt_to_bytes, ConversationKey}; +use nip44::v2::{encrypt_to_bytes, ConversationKey}; use nostr_sdk::prelude::*; -use std::thread::sleep; +use sqlx::SqlitePool; use std::time::Duration; use std::{fs, path::Path}; +use uuid::Uuid; + +#[derive(Clone, Debug)] +pub enum Event { + SmallOrder(SmallOrder), + Dispute(Dispute), // Assuming you have a Dispute struct + MessageTuple(Box<(Message, u64)>), +} + +#[derive(Clone, Debug)] +pub enum ListKind { + Orders, + Disputes, + DirectMessagesUser, + DirectMessagesAdmin, +} + +pub async fn save_order( + order: SmallOrder, + trade_keys: &Keys, + request_id: u64, + trade_index: i64, +) -> Result<()> { + let pool = connect().await?; + if let Ok(order) = Order::new(&pool, order, trade_keys, Some(request_id as i64)).await { + if let Some(order_id) = order.id { + println!("Order {} created", order_id); + } else { + println!("Warning: The newly created order has no ID."); + } + // Update last trade index to be used in next trade + match User::get(&pool).await { + Ok(mut user) => { + user.set_last_trade_index(trade_index); + if let Err(e) = user.save(&pool).await { + println!("Failed to update user: {}", e); + } + } + Err(e) => println!("Failed to get user: {}", e), + } + } + Ok(()) +} + +/// Wait for incoming gift wraps or events coming in +pub async fn wait_for_dm( + client: &Client, + trade_keys: &Keys, + request_id: u64, + trade_index: i64, + mut order: Option, +) -> anyhow::Result<()> { + let mut notifications = client.notifications(); + + match tokio::time::timeout(Duration::from_secs(10), async move { + while let Ok(notification) = notifications.recv().await { + if let RelayPoolNotification::Event { event, .. } = notification { + if event.kind == nostr_sdk::Kind::GiftWrap { + let gift = nip59::extract_rumor(trade_keys, &event).await.unwrap(); + let (message, _): (Message, Option) = serde_json::from_str(&gift.rumor.content).unwrap(); + let message = message.get_inner_message_kind(); + if message.request_id == Some(request_id) { + match message.action { + Action::NewOrder => { + if let Some(Payload::Order(order)) = message.payload.as_ref() { + save_order(order.clone(), trade_keys, request_id, trade_index).await.map_err(|_| ())?; + return Ok(()); + } + } + // this is the case where the buyer adds an invoice to a takesell order + Action::WaitingSellerToPay => { + println!("Now we should wait for the seller to pay the invoice"); + if let Some(mut order) = order.take() { + let pool = connect().await.map_err(|_| ())?; + match order + .set_status(Status::WaitingPayment.to_string()) + .save(&pool) + .await + { + Ok(_) => println!("Order status updated"), + Err(e) => println!("Failed to update order status: {}", e), + } + } + } + // this is the case where the buyer adds an invoice to a takesell order + Action::AddInvoice => { + if let Some(Payload::Order(order)) = &message.payload { + println!( + "Please add a lightning invoice with amount of {}", + order.amount + ); + return Ok(()); + } + } + // this is the case where the buyer pays the invoice coming from a takebuy + Action::PayInvoice => { + if let Some(Payload::PaymentRequest(order, invoice, _)) = &message.payload { + println!( + "Mostro sent you this hold invoice for order id: {}", + order + .as_ref() + .and_then(|o| o.id) + .map_or("unknown".to_string(), |id| id.to_string()) + ); + println!(); + println!("Pay this invoice to continue --> {}", invoice); + println!(); + if let Some(order) = order { + let store_order = order.clone(); + save_order(store_order, trade_keys, request_id, trade_index).await.map_err(|_| ())?; + } + return Ok(()); + } + } + Action::CantDo => { + match message.payload { + Some(Payload::CantDo(Some(CantDoReason::OutOfRangeFiatAmount | CantDoReason::OutOfRangeSatsAmount))) => { + println!("Error: Amount is outside the allowed range. Please check the order's min/max limits."); + return Err(()); + } + Some(Payload::CantDo(Some(CantDoReason::PendingOrderExists))) => { + println!("Error: A pending order already exists. Please wait for it to be filled or canceled."); + return Err(()); + } + Some(Payload::CantDo(Some(CantDoReason::InvalidTradeIndex))) => { + println!("Error: Invalid trade index. Please synchronize the trade index with mostro"); + return Err(()); + } + _ => { + println!("Unknown reason: {:?}", message.payload); + return Err(()); + } + } + } + // this is the case where the user cancels the order + Action::Canceled => { + if let Some(order_id) = &message.id { + // Acquire database connection + let pool = connect().await.map_err(|_| ())?; + // Verify order exists before deletion + if Order::get_by_id(&pool, &order_id.to_string()).await.is_ok() { + Order::delete_by_id(&pool, &order_id.to_string()) + .await + .map_err(|_| ())?; + // Release database connection + drop(pool); + println!("Order {} canceled!", order_id); + return Ok(()); + } else { + println!("Order not found: {}", order_id); + return Err(()); + } + } + } + _ => { + println!("Unknown action: {:?}", message.action); + return Err(()); + } + } + } + } + } + } + Ok(()) + }) + .await { + Ok(_) => Ok(()), + Err(_) => Err(anyhow::anyhow!("Timeout waiting for DM or gift wrap event")) + } +} + +#[derive(Debug, Clone, Copy)] +enum MessageType { + PrivateDirectMessage, + PrivateGiftWrap, + SignedGiftWrap, +} + +fn determine_message_type(to_user: bool, private: bool) -> MessageType { + match (to_user, private) { + (true, _) => MessageType::PrivateDirectMessage, + (false, true) => MessageType::PrivateGiftWrap, + (false, false) => MessageType::SignedGiftWrap, + } +} + +fn create_expiration_tags(expiration: Option) -> Tags { + let mut tags: Vec = Vec::with_capacity(1 + usize::from(expiration.is_some())); + + if let Some(timestamp) = expiration { + tags.push(Tag::expiration(timestamp)); + } + + Tags::from_list(tags) +} + +async fn create_private_dm_event( + trade_keys: &Keys, + receiver_pubkey: &PublicKey, + payload: String, + pow: u8, +) -> Result { + // Derive conversation key + let ck = ConversationKey::derive(trade_keys.secret_key(), receiver_pubkey)?; + // Encrypt payload + let encrypted_content = encrypt_to_bytes(&ck, payload.as_bytes())?; + // Encode with base64 + let b64decoded_content = general_purpose::STANDARD.encode(encrypted_content); + // Compose builder + Ok( + EventBuilder::new(nostr_sdk::Kind::PrivateDirectMessage, b64decoded_content) + .pow(pow) + .tag(Tag::public_key(*receiver_pubkey)) + .sign_with_keys(trade_keys)?, + ) +} + +async fn create_gift_wrap_event( + trade_keys: &Keys, + identity_keys: Option<&Keys>, + receiver_pubkey: &PublicKey, + payload: String, + pow: u8, + expiration: Option, + signed: bool, +) -> Result { + let message = Message::from_json(&payload).unwrap(); + + let content = if signed { + let _identity_keys = identity_keys + .ok_or_else(|| Error::msg("identity_keys required for signed messages"))?; + // We sign the message + let sig = Message::sign(payload, trade_keys); + serde_json::to_string(&(message, sig)).unwrap() + } else { + // We compose the content, when private we don't sign the payload + let content: (Message, Option) = (message, None); + serde_json::to_string(&content).unwrap() + }; + + // We create the rumor + let rumor = EventBuilder::text_note(content) + .pow(pow) + .build(trade_keys.public_key()); + + let tags = create_expiration_tags(expiration); + + let signer_keys = if signed { + identity_keys.unwrap() + } else { + trade_keys + }; + + Ok(EventBuilder::gift_wrap(signer_keys, receiver_pubkey, rumor, tags).await?) +} pub async fn send_dm( client: &Client, @@ -26,60 +282,40 @@ pub async fn send_dm( .unwrap_or("false".to_string()) .parse::() .unwrap(); - let event = if to_user { - // Derive conversation key - let ck = ConversationKey::derive(trade_keys.secret_key(), receiver_pubkey)?; - // Encrypt payload - let encrypted_content = encrypt_to_bytes(&ck, payload.as_bytes())?; - // Encode with base64 - let b64decoded_content = general_purpose::STANDARD.encode(encrypted_content); - // Compose builder - EventBuilder::new(nostr_sdk::Kind::PrivateDirectMessage, b64decoded_content) - .pow(pow) - .tag(Tag::public_key(*receiver_pubkey)) - .sign_with_keys(trade_keys)? - } else if private { - let message = Message::from_json(&payload).unwrap(); - // We compose the content, when private we don't sign the payload - let content: (Message, Option) = (message, None); - let content = serde_json::to_string(&content).unwrap(); - // We create the rumor - let rumor = EventBuilder::text_note(content) - .pow(pow) - .build(trade_keys.public_key()); - let mut tags: Vec = Vec::with_capacity(1 + usize::from(expiration.is_some())); - - if let Some(timestamp) = expiration { - tags.push(Tag::expiration(timestamp)); - } - let tags = Tags::from_list(tags); - EventBuilder::gift_wrap(trade_keys, receiver_pubkey, rumor, tags).await? - } else { - let identity_keys = identity_keys - .ok_or_else(|| Error::msg("identity_keys required when to_user is false"))?; - // We sign the message - let message = Message::from_json(&payload).unwrap(); - let sig = Message::sign(payload.clone(), trade_keys); - // We compose the content - let content = serde_json::to_string(&(message, sig)).unwrap(); - // We create the rumor - let rumor = EventBuilder::text_note(content) - .pow(pow) - .build(trade_keys.public_key()); - let mut tags: Vec = Vec::with_capacity(1 + usize::from(expiration.is_some())); + let message_type = determine_message_type(to_user, private); - if let Some(timestamp) = expiration { - tags.push(Tag::expiration(timestamp)); + let event = match message_type { + MessageType::PrivateDirectMessage => { + create_private_dm_event(trade_keys, receiver_pubkey, payload, pow).await? + } + MessageType::PrivateGiftWrap => { + create_gift_wrap_event( + trade_keys, + identity_keys, + receiver_pubkey, + payload, + pow, + expiration, + false, + ) + .await? + } + MessageType::SignedGiftWrap => { + create_gift_wrap_event( + trade_keys, + identity_keys, + receiver_pubkey, + payload, + pow, + expiration, + true, + ) + .await? } - let tags = Tags::from_list(tags); - - EventBuilder::gift_wrap(identity_keys, receiver_pubkey, rumor, tags).await? }; - info!("Sending event: {event:#?}"); client.send_event(&event).await?; - Ok(()) } @@ -94,6 +330,7 @@ pub async fn connect_nostr() -> Result { for r in relays.into_iter() { client.add_relay(r).await?; } + // Connect to relays and keep connection alive client.connect().await; @@ -106,15 +343,13 @@ pub async fn send_message_sync( trade_keys: &Keys, receiver_pubkey: PublicKey, message: Message, - wait_for_dm: bool, + _wait_for_dm: bool, to_user: bool, -) -> Result> { - let message_json = message.as_json().map_err(|_| Error::msg("Failed to serialize message"))?; - // Send dm to receiver pubkey - println!( - "SENDING DM with trade keys: {:?}", - trade_keys.public_key().to_hex() - ); +) -> Result<()> { + let message_json = message + .as_json() + .map_err(|_| Error::msg("Failed to serialize message"))?; + send_dm( client, identity_keys, @@ -125,271 +360,133 @@ pub async fn send_message_sync( to_user, ) .await?; - // FIXME: This is a hack to wait for the DM to be sent - sleep(Duration::from_secs(2)); - - let dm: Vec<(Message, u64)> = if wait_for_dm { - get_direct_messages(client, trade_keys, 15, to_user).await - } else { - Vec::new() - }; - Ok(dm) + Ok(()) } -pub async fn get_direct_messages( - client: &Client, - my_key: &Keys, - since: i64, - from_user: bool, -) -> Vec<(Message, u64)> { - // We use a fake timestamp to thwart time-analysis attacks - let fake_since = 2880; - let fake_since_time = chrono::Utc::now() - .checked_sub_signed(chrono::Duration::minutes(fake_since)) - .unwrap() - .timestamp() as u64; - - let fake_timestamp = Timestamp::from(fake_since_time); - let filters = if from_user { - let since_time = chrono::Utc::now() - .checked_sub_signed(chrono::Duration::minutes(since)) - .unwrap() - .timestamp() as u64; - let timestamp = Timestamp::from(since_time); - Filter::new() - .kind(nostr_sdk::Kind::PrivateDirectMessage) - .pubkey(my_key.public_key()) - .since(timestamp) - } else { - Filter::new() - .kind(nostr_sdk::Kind::GiftWrap) - .pubkey(my_key.public_key()) - .since(fake_timestamp) - }; - - info!("Request events with event kind : {:?} ", filters.kinds); - - let mut direct_messages: Vec<(Message, u64)> = Vec::new(); - - if let Ok(mostro_req) = client.fetch_events(filters, Duration::from_secs(15)).await { - // Buffer vector for direct messages - // Vector for single order id check - maybe multiple relay could send the same order id? Check unique one... - let mut id_list = Vec::::new(); - - for dm in mostro_req.iter() { - if !id_list.contains(&dm.id) { - id_list.push(dm.id); - let (created_at, message) = if from_user { - let ck = - if let Ok(ck) = ConversationKey::derive(my_key.secret_key(), &dm.pubkey) { - ck - } else { - continue; - }; - let b64decoded_content = - match general_purpose::STANDARD.decode(dm.content.as_bytes()) { - Ok(b64decoded_content) => b64decoded_content, - Err(_) => { - continue; - } - }; - - let unencrypted_content = decrypt_to_bytes(&ck, &b64decoded_content) - .expect("Failed to decrypt message"); - - let message = - String::from_utf8(unencrypted_content).expect("Found invalid UTF-8"); - let message = Message::from_json(&message).expect("Failed on deserializing"); - - (dm.created_at, message) - } else { - let unwrapped_gift = match nip59::extract_rumor(my_key, dm).await { - Ok(u) => u, - Err(_) => { - println!("Error unwrapping gift"); - continue; - } - }; - let (message, _): (Message, Option) = - serde_json::from_str(&unwrapped_gift.rumor.content).unwrap(); - - (unwrapped_gift.rumor.created_at, message) - }; - - // Here we discard messages older than the real since parameter - let since_time = chrono::Utc::now() - .checked_sub_signed(chrono::Duration::minutes(30)) - .unwrap() - .timestamp() as u64; - if created_at.as_u64() < since_time { - continue; - } - direct_messages.push((message, created_at.as_u64())); - } +pub fn create_filter(list_kind: ListKind, pubkey: PublicKey) -> Filter { + match list_kind { + ListKind::Orders => { + let since_time = chrono::Utc::now() + .checked_sub_signed(chrono::Duration::days(7)) + .unwrap() + .timestamp() as u64; + + let timestamp = Timestamp::from(since_time); + + // Create filter for fetching orders + Filter::new() + .author(pubkey) + .limit(50) + .since(timestamp) + .custom_tag(SingleLetterTag::lowercase(Alphabet::Z), "order".to_string()) + .kind(nostr_sdk::Kind::Custom(NOSTR_REPLACEABLE_EVENT_KIND)) + } + ListKind::Disputes => { + let since_time = chrono::Utc::now() + .checked_sub_signed(chrono::Duration::days(7)) + .unwrap() + .timestamp() as u64; + + let timestamp = Timestamp::from(since_time); + + // Create filter for fetching orders + Filter::new() + .author(pubkey) + .limit(50) + .since(timestamp) + .custom_tag( + SingleLetterTag::lowercase(Alphabet::Z), + "dispute".to_string(), + ) + .kind(nostr_sdk::Kind::Custom(NOSTR_REPLACEABLE_EVENT_KIND)) + } + ListKind::DirectMessagesAdmin => { + // We use a fake timestamp to thwart time-analysis attacks + let fake_since = 2880; + let fake_since_time = chrono::Utc::now() + .checked_sub_signed(chrono::Duration::minutes(fake_since)) + .unwrap() + .timestamp() as u64; + + let fake_timestamp = Timestamp::from(fake_since_time); + + Filter::new() + .kind(nostr_sdk::Kind::GiftWrap) + .pubkey(pubkey) + .since(fake_timestamp) + } + ListKind::DirectMessagesUser => { + let since_time = chrono::Utc::now() + .checked_sub_signed(chrono::Duration::minutes(30)) + .unwrap() + .timestamp() as u64; + let timestamp = Timestamp::from(since_time); + Filter::new() + .kind(nostr_sdk::Kind::PrivateDirectMessage) + .pubkey(pubkey) + .since(timestamp) } - // Return element sorted by second tuple element ( Timestamp ) - direct_messages.sort_by(|a, b| a.1.cmp(&b.1)); } - - direct_messages } -pub async fn get_orders_list( - pubkey: PublicKey, - status: Status, +#[allow(clippy::too_many_arguments)] +pub async fn fetch_events_list( + list_kind: ListKind, + status: Option, currency: Option, kind: Option, + mostro_pubkey: PublicKey, + mostro_keys: &Keys, + trade_index: i64, + pool: &SqlitePool, client: &Client, -) -> Result> { - let since_time = chrono::Utc::now() - .checked_sub_signed(chrono::Duration::days(7)) - .unwrap() - .timestamp() as u64; - - let timestamp = Timestamp::from(since_time); - - let filters = Filter::new() - .author(pubkey) - .limit(50) - .since(timestamp) - .custom_tag(SingleLetterTag::lowercase(Alphabet::Z), "order".to_string()) - .kind(nostr_sdk::Kind::Custom(NOSTR_REPLACEABLE_EVENT_KIND)); - - info!( - "Request to mostro id : {:?} with event kind : {:?} ", - filters.authors, filters.kinds - ); - - // Extracted Orders List - let mut complete_events_list = Vec::::new(); - let mut requested_orders_list = Vec::::new(); - - // Send all requests to relays - if let Ok(mostro_req) = client.fetch_events(filters, Duration::from_secs(15)).await { - // Scan events to extract all orders - for el in mostro_req.iter() { - let order = order_from_tags(el.tags.clone()); - - if order.is_err() { - error!("{order:?}"); - continue; - } - let mut order = order?; - - info!("Found Order id : {:?}", order.id.unwrap()); - - if order.id.is_none() { - info!("Order ID is none"); - continue; - } - - if order.kind.is_none() { - info!("Order kind is none"); - continue; - } - - if order.status.is_none() { - info!("Order status is none"); - continue; - } - - // Get created at field from Nostr event - order.created_at = Some(el.created_at.as_u64() as i64); - - complete_events_list.push(order.clone()); - - if order.status.ne(&Some(status)) { - continue; - } - - if currency.is_some() && order.fiat_code.ne(¤cy.clone().unwrap()) { - continue; - } - - if kind.is_some() && order.kind.ne(&kind) { - continue; - } - // Add just requested orders requested by filtering - requested_orders_list.push(order); +) -> Result> { + match list_kind { + ListKind::Orders => { + let filters = create_filter(list_kind, mostro_pubkey); + let fetched_events = client + .fetch_events(filters, Duration::from_secs(15)) + .await?; + let orders = parse_orders_events(fetched_events, currency, status, kind); + Ok(orders.into_iter().map(Event::SmallOrder).collect()) } - } - - // Order all element ( orders ) received to filter - discard disaligned messages - // if an order has an older message with the state we received is discarded for the latest one - requested_orders_list.retain(|keep| { - !complete_events_list - .iter() - .any(|x| x.id == keep.id && x.created_at > keep.created_at) - }); - // Sort by id to remove duplicates - requested_orders_list.sort_by(|a, b| b.id.cmp(&a.id)); - requested_orders_list.dedup_by(|a, b| a.id == b.id); - - // Finally sort list by creation time - requested_orders_list.sort_by(|a, b| b.created_at.cmp(&a.created_at)); - - Ok(requested_orders_list) -} - -pub async fn get_disputes_list(pubkey: PublicKey, client: &Client) -> Result> { - let since_time = chrono::Utc::now() - .checked_sub_signed(chrono::Duration::days(7)) - .unwrap() - .timestamp() as u64; - - let timestamp = Timestamp::from(since_time); - - let filter = Filter::new() - .author(pubkey) - .limit(50) - .since(timestamp) - .custom_tag( - SingleLetterTag::lowercase(Alphabet::Z), - "dispute".to_string(), - ) - .kind(nostr_sdk::Kind::Custom(NOSTR_REPLACEABLE_EVENT_KIND)); - - // Extracted Orders List - let mut disputes_list = Vec::::new(); - - // Send all requests to relays - if let Ok(mostro_req) = client.fetch_events(filter, Duration::from_secs(15)).await { - // Scan events to extract all disputes - for d in mostro_req.iter() { - let dispute = dispute_from_tags(d.tags.clone()); - - if dispute.is_err() { - error!("{dispute:?}"); - continue; + ListKind::DirectMessagesAdmin => { + let filters = create_filter(list_kind, mostro_keys.public_key()); + let fetched_events = client + .fetch_events(filters, Duration::from_secs(15)) + .await?; + let direct_messages_mostro = parse_dm_events(fetched_events, mostro_keys).await; + Ok(direct_messages_mostro + .into_iter() + .map(|t| Event::MessageTuple(Box::new(t))) + .collect()) + } + ListKind::DirectMessagesUser => { + let mut direct_messages: Vec<(Message, u64)> = Vec::new(); + for index in 1..=trade_index { + let trade_key = User::get_trade_keys(pool, index).await?; + let filter = create_filter(ListKind::DirectMessagesUser, trade_key.public_key()); + let fetched_user_messages = + client.fetch_events(filter, Duration::from_secs(15)).await?; + let direct_messages_for_trade_key = + parse_dm_events(fetched_user_messages, &trade_key).await; + direct_messages.extend(direct_messages_for_trade_key); } - let mut dispute = dispute?; - - info!("Found Dispute id : {:?}", dispute.id); - - // Get created at field from Nostr event - dispute.created_at = d.created_at.as_u64() as i64; - disputes_list.push(dispute); + Ok(direct_messages + .into_iter() + .map(|t| Event::MessageTuple(Box::new(t))) + .collect()) + } + ListKind::Disputes => { + let filters = create_filter(list_kind, mostro_pubkey); + let fetched_events = client + .fetch_events(filters, Duration::from_secs(15)) + .await?; + let disputes = parse_dispute_events(fetched_events); + Ok(disputes.into_iter().map(Event::Dispute).collect()) } } - - let buffer_dispute_list = disputes_list.clone(); - // Order all element ( orders ) received to filter - discard disaligned messages - // if an order has an older message with the state we received is discarded for the latest one - disputes_list.retain(|keep| { - !buffer_dispute_list - .iter() - .any(|x| x.id == keep.id && x.created_at > keep.created_at) - }); - - // Sort by id to remove duplicates - disputes_list.sort_by(|a, b| b.id.cmp(&a.id)); - disputes_list.dedup_by(|a, b| a.id == b.id); - - // Finally sort list by creation time - disputes_list.sort_by(|a, b| b.created_at.cmp(&a.created_at)); - - Ok(disputes_list) } /// Uppercase first letter of a string. @@ -411,3 +508,21 @@ pub fn get_mcli_path() -> String { mcli_path } + +pub async fn run_simple_order_msg( + command: Commands, + order_id: &Uuid, + identity_keys: &Keys, + mostro_key: PublicKey, + client: &Client, +) -> Result<()> { + execute_send_msg( + command, + Some(*order_id), + Some(identity_keys), + mostro_key, + client, + None, + ) + .await +}