Skip to content
Open
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
131 changes: 129 additions & 2 deletions rust/src/api/orders.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1427,10 +1427,14 @@ pub(crate) async fn subscribe_gift_wraps(trade_pubkey: nostr_sdk::PublicKey, tra
.author(mostro_pubkey)
.pubkey(trade_pubkey)
.limit(0);
if let Err(e) = client.subscribe(filter, None).await {
let sub_id = trade_subscription_id(&trade_pubkey);
if let Err(e) = client.subscribe_with_id(sub_id.clone(), filter, None).await {
log::warn!("[orders] subscribe_gift_wraps subscribe failed: {e}");
return;
}
// Claim ownership of this (deterministic) id. A later re-subscribe for the
// same trade supersedes us; on exit we only unsubscribe if still current.
let sub_gen = claim_subscription(&sub_id);

let trade_pubkey_hex = trade_pubkey.to_hex();
crate::api::logging::blog_info("orders", format!(
Expand Down Expand Up @@ -1517,6 +1521,17 @@ pub(crate) async fn subscribe_gift_wraps(trade_pubkey: nostr_sdk::PublicKey, tra
// (request timed out and no genuine late reply ever arrived) is dead
// state — drop it, whatever attempt it belongs to.
purge_pending_request(&trade_pubkey_hex);
// Tear down the relay-side subscription on every exit path (idle
// timeout, shutdown, closed) so it never outlives the task — but only
// if we still own the id. A newer watcher may have re-subscribed under
// the same deterministic id (retry within the idle window); unsubscribing
// then would kill *its* live subscription. Re-fetch the pool rather than
// hold `client` across the loop (same approach as subscribe_incoming_chat).
if owns_subscription(&sub_id, sub_gen) {
if let Ok(pool) = crate::api::nostr::get_pool() {
pool.client().unsubscribe(&sub_id).await;
}
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
});
}

Expand Down Expand Up @@ -2350,10 +2365,14 @@ async fn subscribe_single_order(order_id: &str) {

let mut rx = client.notifications();
let filter = crate::nostr::order_events::trade_order_filter(&mostro_pubkey, &order_id);
if let Err(e) = client.subscribe(filter, None).await {
let sub_id = single_order_subscription_id(&order_id);
if let Err(e) = client.subscribe_with_id(sub_id.clone(), filter, None).await {
log::warn!("[orders] subscribe_single_order subscribe failed: {e}");
return;
}
// Claim ownership of this (deterministic) id; only the current owner
// unsubscribes on exit (see subscribe_gift_wraps for the rationale).
let sub_gen = claim_subscription(&sub_id);
log::info!("[orders] subscribed to d-tag updates for order={order_id}");

use nostr_sdk::RelayPoolNotification;
Expand Down Expand Up @@ -2410,6 +2429,14 @@ async fn subscribe_single_order(order_id: &str) {
Ok(Ok(_)) => continue,
}
}
// Tear down the relay-side subscription on every exit path (idle
// timeout, shutdown, closed) so the single-order watch never outlives
// the task — but only if we still own the id, so a stale watcher never
// unsubscribes a newer one's live subscription. `client` is already
// owned by this closure.
if owns_subscription(&sub_id, sub_gen) {
client.unsubscribe(&sub_id).await;
}
});
}

Expand Down Expand Up @@ -2532,6 +2559,57 @@ fn orders_subscription_id() -> nostr_sdk::SubscriptionId {
nostr_sdk::SubscriptionId::new("mostro-orders")
}

/// Stable per-trade subscription ID for the ephemeral Kind 14 Mostro-reply
/// feed created by `subscribe_gift_wraps`. Deterministic (full trade pubkey
/// hex, not a prefix — no collision) so a repeat subscribe for the same trade
/// key replaces in place instead of stacking a second relay subscription.
fn trade_subscription_id(trade_pubkey: &nostr_sdk::PublicKey) -> nostr_sdk::SubscriptionId {
nostr_sdk::SubscriptionId::new(format!("mostro-trade-{}", trade_pubkey.to_hex()))
}

/// Stable subscription ID for a single-order (`d`-tag) Kind 38383 watch
/// created by `subscribe_single_order`.
fn single_order_subscription_id(order_id: &str) -> nostr_sdk::SubscriptionId {
nostr_sdk::SubscriptionId::new(format!("mostro-order-{order_id}"))
}

/// Per-subscription-id generation counter.
///
/// A deterministic subscription id is reused across re-subscribes for the same
/// trade/order (e.g. a retry within the 30-min idle window). Without ownership
/// tracking, an earlier watcher's exit-path `unsubscribe` would tear down the
/// *replacement* watcher's live subscription on that same id. Each watcher
/// claims the id on subscribe (bumping the generation) and only unsubscribes on
/// exit if it still owns the current generation — a stale watcher skips cleanup
/// and lets the newer owner keep the subscription.
static SUBSCRIPTION_GENERATIONS: OnceLock<std::sync::Mutex<HashMap<nostr_sdk::SubscriptionId, u64>>> =
OnceLock::new();

fn subscription_generations() -> &'static std::sync::Mutex<HashMap<nostr_sdk::SubscriptionId, u64>> {
SUBSCRIPTION_GENERATIONS.get_or_init(|| std::sync::Mutex::new(HashMap::new()))
}

/// Claim `id` for the calling watcher: bump its generation and return the new
/// value. A later watcher claiming the same id bumps it again, so this caller
/// can later detect it has been superseded.
fn claim_subscription(id: &nostr_sdk::SubscriptionId) -> u64 {
let mut gens = subscription_generations().lock().unwrap();
let g = gens.entry(id.clone()).or_insert(0);
*g += 1;
*g
}

/// True only if `my_gen` is still the current generation for `id` — i.e. no
/// newer watcher has claimed it. A stale watcher must NOT unsubscribe, or it
/// would kill the replacement's live subscription.
fn owns_subscription(id: &nostr_sdk::SubscriptionId, my_gen: u64) -> bool {
subscription_generations()
.lock()
.unwrap()
.get(id)
.is_some_and(|current| *current == my_gen)
}

/// Stable subscription ID for the Kind 14 Mostro-reply feed.
fn mostro_dm_subscription_id() -> nostr_sdk::SubscriptionId {
nostr_sdk::SubscriptionId::new("mostro-dm")
Expand Down Expand Up @@ -3198,6 +3276,55 @@ mod tests {
assert_eq!(admin_pubkey_from_payload(Some(&Payload::Amount(42))), None);
}

// ── deterministic subscription ids (#182) ─────────────────────────────────
#[test]
fn trade_subscription_id_is_deterministic_and_per_pubkey() {
let pk = nostr_sdk::Keys::generate().public_key();
// Idempotent: the same trade key always maps to the same id, so a
// repeat subscribe replaces in place instead of stacking a new
// relay-side subscription.
assert_eq!(trade_subscription_id(&pk), trade_subscription_id(&pk));
// Full pubkey hex (not a prefix) — no cross-trade collision.
assert_eq!(
trade_subscription_id(&pk),
nostr_sdk::SubscriptionId::new(format!("mostro-trade-{}", pk.to_hex())),
);
let other = nostr_sdk::Keys::generate().public_key();
assert_ne!(trade_subscription_id(&other), trade_subscription_id(&pk));
}

#[test]
fn single_order_subscription_id_matches_expected_format() {
assert_eq!(
single_order_subscription_id("abc123"),
nostr_sdk::SubscriptionId::new("mostro-order-abc123"),
);
}

Comment thread
coderabbitai[bot] marked this conversation as resolved.
#[test]
fn a_stale_watcher_does_not_unsubscribe_a_replacement() {
// Two watchers claim the same deterministic id in turn (a re-subscribe
// for the same trade — e.g. a retry within the idle window).
let id = nostr_sdk::SubscriptionId::new("mostro-trade-lifecycle-test");
let gen_a = claim_subscription(&id);
let gen_b = claim_subscription(&id);
assert_ne!(gen_a, gen_b, "each claim must advance the generation");
// Watcher A is now stale: it must NOT unsubscribe on exit, or it would
// tear down watcher B's live subscription on the shared id.
assert!(
!owns_subscription(&id, gen_a),
"the superseded watcher must not own the id",
);
// Watcher B still owns the id, so it (and only it) unsubscribes on exit.
assert!(
owns_subscription(&id, gen_b),
"the current watcher must own the id",
);
// Distinct ids are tracked independently.
let other = nostr_sdk::SubscriptionId::new("mostro-trade-other");
assert!(owns_subscription(&other, claim_subscription(&other)));
}

fn insert_pending_create(key: &str, request_id: u64) -> tokio::sync::oneshot::Receiver<DaemonReply> {
let (tx, rx) = tokio::sync::oneshot::channel::<DaemonReply>();
pending_requests().lock().unwrap().insert(
Expand Down
Loading