Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
87 changes: 77 additions & 10 deletions src/app/restore_session.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
use crate::app::context::AppContext;
use crate::{db::RestoreSessionManager, util::enqueue_restore_session_msg};
use crate::config::settings::get_db_pool;
use crate::db::{find_failed_payment_for_master_key, RestoreSessionManager};
use crate::util::{enqueue_order_msg_on_restore_queue, enqueue_restore_session_msg};
use mostro_core::prelude::*;
use nostr_sdk::prelude::*;

Expand Down Expand Up @@ -42,14 +44,18 @@ pub async fn restore_session_action(

// Start a background task to handle the results
tokio::spawn(async move {
handle_restore_session_results(manager, trade_key).await;
handle_restore_session_results(manager, trade_key, master_key).await;
});

Ok(())
}

/// Handle restore session results in the background
async fn handle_restore_session_results(mut manager: RestoreSessionManager, trade_key: String) {
async fn handle_restore_session_results(
mut manager: RestoreSessionManager,
trade_key: String,
master_key: String,
) {
// Wait for the result with a timeout
let timeout = tokio::time::Duration::from_secs(60 * 60); // 1 hour timeout

Expand All @@ -58,6 +64,7 @@ async fn handle_restore_session_results(mut manager: RestoreSessionManager, trad
// Send the restore session response
if let Err(e) = send_restore_session_response(
&trade_key,
&master_key,
result.restore_orders,
result.restore_disputes,
)
Expand All @@ -82,6 +89,7 @@ async fn handle_restore_session_results(mut manager: RestoreSessionManager, trad
/// Send restore session response to the user
async fn send_restore_session_response(
trade_key: &str,
master_key: &str,
orders: Vec<RestoredOrdersInfo>,
disputes: Vec<RestoredDisputesInfo>,
) -> Result<(), MostroError> {
Expand All @@ -102,6 +110,58 @@ async fn send_restore_session_response(
// No key in the log line (AGENTS.md: scrub Nostr keys from logs).
tracing::info!("Restore session response sent to user");

// Re-send AddInvoice for any orders stuck in settled-hold-invoice with a failed payment.
// Uses get_db_pool() since pool is not available in this function's scope.
//
// `find_failed_payment_for_master_key` filters on `master_buyer_pubkey =
// master_key`, so every order returned belongs to this restoring user as
// the BUYER. The correct AddInvoice recipient is therefore the order's
// own buyer trade key (`order.buyer_pubkey` / get_buyer_pubkey()) — the
// key the client actually listens on for order DMs. Sending to the
// master/identity key instead would (a) never reach the client, since
// order messages are only delivered to trade keys, and (b) publish the
// identity key as a gift-wrap recipient on Nostr, linking it to this
// order and breaking trade-key unlinkability.
let pool = get_db_pool();
match find_failed_payment_for_master_key(&pool, master_key).await {
Ok(failed_orders) => {
for order in failed_orders {
let buyer_trade_pubkey = match order.get_buyer_pubkey() {
Ok(pk) => pk,
Err(_) => {
tracing::warn!(
"Order {} has no valid buyer_pubkey (trade key); skipping AddInvoice re-send",
order.id
);
continue;
}
};
// Route through the restore-session queue so this AddInvoice
// is sent AFTER the restore-session response (same queue =
// FIFO), not before it.
enqueue_order_msg_on_restore_queue(
Some(order.id),
Action::AddInvoice,
Some(Payload::Order(SmallOrder::from(order.clone()))),
buyer_trade_pubkey,
order.trade_index_buyer,
)
.await;
// No key in the log line (AGENTS.md: scrub Nostr keys from logs).
tracing::info!(
"Re-sent AddInvoice for order {} to buyer trade key (failed payment)",
order.id
);
}
}
Err(e) => {
tracing::error!(
"Failed to query failed payments during restore-session: {}",
e
);
}
}

Ok(())
}

Expand Down Expand Up @@ -191,11 +251,11 @@ mod tests {

let manager = RestoreSessionManager::new();
manager
.start_restore_session(pool.clone(), master_key)
.start_restore_session(pool.clone(), master_key.clone())
.await
.unwrap();

handle_restore_session_results(manager, trade_key).await;
handle_restore_session_results(manager, trade_key, master_key).await;

let queued = queued_restore_msgs_for(&trade_pubkey).await;
assert_eq!(queued.len(), 1);
Expand All @@ -214,20 +274,25 @@ mod tests {

let manager = RestoreSessionManager::new();
manager
.start_restore_session(pool.clone(), master_key)
.start_restore_session(pool.clone(), master_key.clone())
.await
.unwrap();

// Must not panic; the send failure is logged and swallowed
handle_restore_session_results(manager, "not-a-hex-key".to_string()).await;
handle_restore_session_results(manager, "not-a-hex-key".to_string(), master_key).await;
}

#[tokio::test]
async fn send_restore_session_response_queues_message_for_valid_key() {
let trade_pubkey = Keys::generate().public_key();

let result =
send_restore_session_response(&trade_pubkey.to_string(), Vec::new(), Vec::new()).await;
let result = send_restore_session_response(
&trade_pubkey.to_string(),
&Keys::generate().public_key().to_string(),
Vec::new(),
Vec::new(),
)
.await;

assert!(result.is_ok());
let queued = queued_restore_msgs_for(&trade_pubkey).await;
Expand All @@ -240,7 +305,9 @@ mod tests {

#[tokio::test]
async fn send_restore_session_response_rejects_invalid_key() {
let result = send_restore_session_response("invalid-key", Vec::new(), Vec::new()).await;
let master_key = Keys::generate().public_key().to_string();
let result =
send_restore_session_response("invalid-key", &master_key, Vec::new(), Vec::new()).await;

assert!(matches!(
result,
Expand Down
134 changes: 134 additions & 0 deletions src/db.rs
Original file line number Diff line number Diff line change
Expand Up @@ -627,6 +627,33 @@ pub async fn edit_pubkeys_order(pool: &SqlitePool, order: &Order) -> Result<Orde
Ok(order)
}

/// Returns orders with `failed_payment = true` and `status = "settled-hold-invoice"`
/// for the given buyer's `master_buyer_pubkey`. Used during restore-session to
/// re-send `Action::AddInvoice` to buyers whose Lightning payment failed.
pub async fn find_failed_payment_for_master_key(
pool: &SqlitePool,
master_key: &str,
) -> Result<Vec<Order>, MostroError> {
// Validate public key format (32-bytes hex)
if !master_key.chars().all(|c| c.is_ascii_hexdigit()) || master_key.len() != 64 {
return Err(MostroCantDo(CantDoReason::InvalidPubkey));
}
let orders = sqlx::query_as::<_, Order>(
r#"
SELECT *
FROM orders
WHERE failed_payment = true
AND status = 'settled-hold-invoice'
AND master_buyer_pubkey = ?1
"#,
)
.bind(master_key)
.fetch_all(pool)
.await
.map_err(|e| MostroInternalErr(ServiceError::DbAccessError(e.to_string())))?;
Ok(orders)
}

pub async fn find_order_by_hash(pool: &SqlitePool, hash: &str) -> Result<Order, MostroError> {
let order = sqlx::query_as::<_, Order>(
r#"
Expand Down Expand Up @@ -2085,6 +2112,113 @@ mod tests {
);
}

// -- Tests for find_failed_payment_for_master_key --
#[tokio::test]
async fn test_find_failed_payment_for_master_key_returns_matching() {
let pool = setup_orders_db().await.unwrap();
let master_key = "a".repeat(64);
sqlx::query(
r#"INSERT INTO orders (id, kind, event_id, status, premium, payment_method,
amount, fiat_code, fiat_amount, created_at, expires_at,
failed_payment, payment_attempts, dev_fee, dev_fee_paid,
master_buyer_pubkey)
VALUES (?1, 'buy', 'ev1', 'settled-hold-invoice', 0, 'lightning',
100000, 'USD', 100, 1700000000, 1700086400,
1, 3, 0, 0, ?2)"#,
)
.bind(uuid::Uuid::new_v4())
.bind(&master_key)
.execute(&pool)
.await
.unwrap();
let result = super::find_failed_payment_for_master_key(&pool, &master_key)
.await
.unwrap();
assert_eq!(result.len(), 1, "Should find matching failed payment order");
}

#[tokio::test]
async fn test_find_failed_payment_for_master_key_ignores_different_key() {
let pool = setup_orders_db().await.unwrap();
let master_key = "a".repeat(64);
let other_key = "b".repeat(64);
sqlx::query(
r#"INSERT INTO orders (id, kind, event_id, status, premium, payment_method,
amount, fiat_code, fiat_amount, created_at, expires_at,
failed_payment, payment_attempts, dev_fee, dev_fee_paid,
master_buyer_pubkey)
VALUES (?1, 'buy', 'ev1', 'settled-hold-invoice', 0, 'lightning',
100000, 'USD', 100, 1700000000, 1700086400,
1, 3, 0, 0, ?2)"#,
)
.bind(uuid::Uuid::new_v4())
.bind(&other_key)
.execute(&pool)
.await
.unwrap();
let result = super::find_failed_payment_for_master_key(&pool, &master_key)
.await
.unwrap();
assert!(
result.is_empty(),
"Should not return orders for a different master key"
);
}

#[tokio::test]
async fn test_find_failed_payment_for_master_key_ignores_non_failed() {
let pool = setup_orders_db().await.unwrap();
let master_key = "a".repeat(64);
sqlx::query(
r#"INSERT INTO orders (id, kind, event_id, status, premium, payment_method,
amount, fiat_code, fiat_amount, created_at, expires_at,
failed_payment, payment_attempts, dev_fee, dev_fee_paid,
master_buyer_pubkey)
VALUES (?1, 'buy', 'ev1', 'settled-hold-invoice', 0, 'lightning',
100000, 'USD', 100, 1700000000, 1700086400,
0, 0, 0, 0, ?2)"#,
)
.bind(uuid::Uuid::new_v4())
.bind(&master_key)
.execute(&pool)
.await
.unwrap();
let result = super::find_failed_payment_for_master_key(&pool, &master_key)
.await
.unwrap();
assert!(
result.is_empty(),
"Should not return orders where failed_payment is false"
);
}

#[tokio::test]
async fn test_find_failed_payment_for_master_key_ignores_wrong_status() {
let pool = setup_orders_db().await.unwrap();
let master_key = "a".repeat(64);
sqlx::query(
r#"INSERT INTO orders (id, kind, event_id, status, premium, payment_method,
amount, fiat_code, fiat_amount, created_at, expires_at,
failed_payment, payment_attempts, dev_fee, dev_fee_paid,
master_buyer_pubkey)
VALUES (?1, 'buy', 'ev1', 'active', 0, 'lightning',
100000, 'USD', 100, 1700000000, 1700086400,
1, 3, 0, 0, ?2)"#,
)
.bind(uuid::Uuid::new_v4())
.bind(&master_key)
.execute(&pool)
.await
.unwrap();
let result = super::find_failed_payment_for_master_key(&pool, &master_key)
.await
.unwrap();
assert!(
result.is_empty(),
"Should not return orders with wrong status"
);
}

// -- Tests for find_order_by_hash --

#[tokio::test]
Expand Down
24 changes: 24 additions & 0 deletions src/util.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1299,6 +1299,30 @@ pub async fn enqueue_order_msg(
.push((message, destination_key));
}

/// Enqueue an order-type message onto the restore-session queue.
///
/// The scheduler drains `queue_order_msg` before `queue_restore_session_msg`
/// in each tick, so an AddInvoice enqueued via `enqueue_order_msg` would be
/// sent BEFORE the restore-session response even when enqueued after it.
/// Routing the AddInvoice through the restore-session queue instead preserves
/// FIFO ordering relative to the restore response: the client receives the
/// restore-session list first, then the AddInvoice prompt for an order it now
/// knows about.
pub async fn enqueue_order_msg_on_restore_queue(
order_id: Option<Uuid>,
action: Action,
payload: Option<Payload>,
destination_key: PublicKey,
trade_index: Option<i64>,
) {
let message = Message::new_order(order_id, None, trade_index, action, payload);
MESSAGE_QUEUES
.queue_restore_session_msg
.write()
.await
.push((message, destination_key));
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}

pub fn get_fiat_amount_requested(order: &Order, msg: &Message) -> Option<i64> {
// Check if order is range and get amount request after checking boundaries
// set order fiat amount to the value requested preparing for hold invoice
Expand Down