Skip to content
Merged
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
4 changes: 2 additions & 2 deletions lib/src/rust/api/messages.dart

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

10 changes: 5 additions & 5 deletions lib/src/rust/api/push.dart

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

121 changes: 97 additions & 24 deletions rust/src/api/messages.rs
Original file line number Diff line number Diff line change
Expand Up @@ -315,7 +315,23 @@ pub(crate) async fn publish_chat_payload_for(
ctx: &ChatContext,
payload: &str,
) -> Result<nostr_sdk::prelude::Event> {
publish_chat_payload(ctx, payload).await
publish_chat_payload(ctx, payload).await.map(|p| p.inner)
}

/// A chat envelope handed to the pool.
struct PublishedChat {
/// The signed inner event: the message's durable identity.
inner: nostr_sdk::prelude::Event,
/// Whether at least one relay accepted the envelope. `send_event` returns
/// `Ok` even when every relay rejected it.
delivered: bool,
}

/// The peer to wake after a publish, if any: only an envelope some relay
/// accepted is worth a wake. Waking for one that reached nobody rings the
/// peer for nothing and debounces the wake of the retry that does land.
fn peer_to_wake(delivered: bool, peer_hex: &str) -> Option<&str> {
delivered.then_some(peer_hex)
}

/// Record a message we just sent to the solver, mirroring what `send_message`
Expand All @@ -342,7 +358,7 @@ pub(crate) async fn store_outgoing_admin_message(
let _ = message_store().add_message(msg).await;
}

async fn publish_chat_payload(ctx: &ChatContext, payload: &str) -> Result<nostr_sdk::prelude::Event> {
async fn publish_chat_payload(ctx: &ChatContext, payload: &str) -> Result<PublishedChat> {
let (outer, inner) =
crate::nostr::transport::mostro_wrap(&ctx.trade_keys, &ctx.conv, &ctx.sign, payload)
.await?;
Expand Down Expand Up @@ -375,7 +391,10 @@ async fn publish_chat_payload(ctx: &ChatContext, payload: &str) -> Result<nostr_
),
);
}
Ok(inner)
Ok(PublishedChat {
inner,
delivered: !output.success.is_empty(),
})
}

/// Send an encrypted text message to the trade counterparty.
Expand Down Expand Up @@ -433,9 +452,15 @@ pub async fn send_message(trade_id: String, content: String) -> Result<ChatMessa
return Err(e);
}
Err(e) => log::warn!("[messages] send_message trade={trade_id}: {e}"),
Ok(inner) => {
id = inner.id.to_hex();
created_at = inner.created_at.as_secs() as i64;
Ok(published) => {
id = published.inner.id.to_hex();
created_at = published.inner.created_at.as_secs() as i64;
// The envelope's `p` tag is `pub(K_conv)`, which
// the push server cannot match: ask it to ring the
// peer's trade key (docs/PUSH_NOTIFICATIONS.md §7.3).
if let Some(peer) = peer_to_wake(published.delivered, peer_hex) {
crate::api::push::wake_peer(peer);
}
}
}
}
Expand Down Expand Up @@ -570,9 +595,12 @@ pub async fn send_file(
Err(e) => log::warn!("[messages] send_file trade={trade_id}: {e}"),
Ok(ctx) => match publish_chat_payload(&ctx, &payload).await {
Err(e) => log::warn!("[messages] send_file trade={trade_id}: {e}"),
Ok(inner) => {
published_id = Some(inner.id.to_hex());
msg_created_at = inner.created_at.as_secs() as i64;
Ok(published) => {
published_id = Some(published.inner.id.to_hex());
msg_created_at = published.inner.created_at.as_secs() as i64;
if let Some(peer) = peer_to_wake(published.delivered, peer_hex) {
crate::api::push::wake_peer(peer);
}
}
},
}
Expand Down Expand Up @@ -834,7 +862,9 @@ async fn persist_decrypted_attachment(
_file_name: &str,
_data: &[u8],
) -> Result<String> {
Err(anyhow!("attachment download to disk is not supported on web"))
Err(anyhow!(
"attachment download to disk is not supported on web"
))
}

fn is_supported_mime_type(mime: &str) -> bool {
Expand Down Expand Up @@ -883,8 +913,14 @@ const MAX_STORED_BYTES_PER_TRADE: usize = 5 * 1024 * 1024;
/// Subscription id for the chat envelope of one order — explicit so every
/// exit path can unsubscribe and a lingering relay subscription never
/// outlives its task.
fn chat_subscription_id(channel: ChatChannel, order_id: &str) -> nostr_sdk::prelude::SubscriptionId {
nostr_sdk::prelude::SubscriptionId::new(format!("mostro-chat-{}{order_id}", channel.id_prefix()))
fn chat_subscription_id(
channel: ChatChannel,
order_id: &str,
) -> nostr_sdk::prelude::SubscriptionId {
nostr_sdk::prelude::SubscriptionId::new(format!(
"mostro-chat-{}{order_id}",
channel.id_prefix()
))
}

/// Which conversation an envelope subscription serves.
Expand Down Expand Up @@ -935,7 +971,6 @@ impl ChatChannel {
fn cursor_key(self, order_id: &str) -> String {
crate::db::settings_keys::chat_cursor(&format!("{}{order_id}", self.id_prefix()))
}

}

/// Orders with a live chat task. Single-owner guard: the peer-reveal capture
Expand Down Expand Up @@ -1110,7 +1145,10 @@ pub(crate) async fn subscribe_incoming_chat(

// Cleanup on every exit path: release ownership and drop the relay
// subscriptions so they never outlive the task.
active_chats().lock().await.remove(&channel.guard_key(&order_id));
active_chats()
.lock()
.await
.remove(&channel.guard_key(&order_id));
if let Ok(pool) = crate::api::nostr::get_pool() {
let client = pool.client();
// Unsubscribing a subscription that already went away is not an
Expand Down Expand Up @@ -1253,8 +1291,17 @@ async fn run_chat_subscription(
..
}) => {
if subscription_id == sub_id {
handle_chat_event(channel, order_id, &allowed_signers, conv, &sign_pubkey, &my_trade_pubkey, &event, &mut state)
.await;
handle_chat_event(
channel,
order_id,
&allowed_signers,
conv,
&sign_pubkey,
&my_trade_pubkey,
&event,
&mut state,
)
.await;
}
if state.flooded {
return;
Expand Down Expand Up @@ -1427,7 +1474,12 @@ pub(crate) async fn resubscribe_active_chats() {
};
log::info!("[messages] resubscribing chat order={order_id}");
crate::rt::spawn(subscribe_incoming_chat(
ChatChannel::Peer, order_id, trade_keys, peer, conv, sign,
ChatChannel::Peer,
order_id,
trade_keys,
peer,
conv,
sign,
));
}
}
Expand Down Expand Up @@ -1561,6 +1613,20 @@ async fn rebuild_session(
mod tests {
use super::*;

const PEER: &str = "aa11aa11aa11aa11aa11aa11aa11aa11aa11aa11aa11aa11aa11aa11aa11aa11";

#[test]
fn a_message_no_relay_accepted_wakes_nobody() {
// `send_message` and `send_file` both gate `wake_peer` on this: an
// empty `output.success` still comes back as `Ok` from the pool.
assert_eq!(peer_to_wake(false, PEER), None);
}

#[test]
fn a_message_some_relay_accepted_wakes_the_peer() {
assert_eq!(peer_to_wake(true, PEER), Some(PEER));
}

#[test]
fn the_two_channels_of_one_order_never_collide() {
let order = "order-1";
Expand Down Expand Up @@ -1924,7 +1990,10 @@ mod tests {
// misinterpreted.
let (content, att) = parse_chat_payload(r#"{"type":"file","url":"x"}"#);
assert_eq!(content, r#"{"type":"file","url":"x"}"#);
assert!(att.is_none(), "incomplete pointer must not become an attachment");
assert!(
att.is_none(),
"incomplete pointer must not become an attachment"
);
}

#[test]
Expand Down Expand Up @@ -2145,7 +2214,10 @@ mod tests {
let trade_keys = nostr_sdk::prelude::Keys::generate();
let trade = live_trade(&order_id, "not-a-pubkey", 1);
assert!(rebuild_session(&trade, &trade_keys).await.is_none());
assert!(crate::mostro::session::session_manager().get_session(&order_id).await.is_none());
assert!(crate::mostro::session::session_manager()
.get_session(&order_id)
.await
.is_none());
}

/// The seam test for #381: `session_or_rebuild` is only useful if the
Expand All @@ -2163,10 +2235,8 @@ mod tests {
#[tokio::test]
#[ignore = "claims the process-global app_db and identity — run with --ignored"]
async fn send_message_rebuilds_session_from_trade_row() {
let db_path = std::env::temp_dir().join(format!(
"mostro-381-seam-test-{}.db",
uuid::Uuid::new_v4()
));
let db_path =
std::env::temp_dir().join(format!("mostro-381-seam-test-{}.db", uuid::Uuid::new_v4()));
crate::db::app_db::init_db(db_path.to_str().unwrap())
.await
.expect("init app db");
Expand Down Expand Up @@ -2195,7 +2265,10 @@ mod tests {
.await
.expect("save the post-reveal row");
assert!(
crate::mostro::session::session_manager().get_session(&order_id).await.is_none(),
crate::mostro::session::session_manager()
.get_session(&order_id)
.await
.is_none(),
"the restart shape: row persisted, session gone"
);

Expand Down
Loading
Loading