1use nostr_sdk::prelude::*;
10
11use crate::rumor::{RumorProcessingResult, RumorEvent, RumorContext, ConversationType, process_rumor};
12use crate::types::Message;
13use crate::state::WRAPPER_ID_CACHE;
14
15pub trait InboundEventHandler: Send + Sync {
20 fn on_dm_received(&self, _chat_id: &str, _msg: &Message, _is_new: bool) {}
22
23 fn on_file_received(&self, _chat_id: &str, _msg: &Message, _is_new: bool) {}
25
26 fn on_reaction_received(&self, _chat_id: &str, _msg: &Message) {}
28
29 fn on_message_deleted(&self, _chat_id: &str, _message_id: &str) {}
32
33 fn on_community_invite(&self, _community_id: &str) {}
37
38 fn on_community_message(&self, _chat_id: &str, _msg: &Message, _is_new: bool) {}
42
43 fn on_community_update(&self, _chat_id: &str, _target_id: &str, _msg: &Message) {}
46
47 fn on_community_removed(&self, _chat_id: &str, _target_id: &str) {}
49
50 #[allow(clippy::too_many_arguments)]
53 fn on_community_presence(
54 &self,
55 _chat_id: &str,
56 _npub: &str,
57 _joined: bool,
58 _event_id: &str,
59 _created_at: u64,
60 _invited_by: Option<&str>,
61 _invited_label: Option<&str>,
62 ) {}
63
64 fn on_community_typing(&self, _chat_id: &str, _npub: &str, _until: u64) {}
66
67 #[allow(clippy::too_many_arguments)]
69 fn on_community_webxdc(
70 &self,
71 _chat_id: &str,
72 _npub: &str,
73 _topic_id: &str,
74 _node_addr: Option<&str>,
75 _event_id: &str,
76 _created_at: u64,
77 ) {}
78
79 fn on_community_self_removed(&self, _community_id: &str) {}
82
83 fn on_community_refreshed(&self, _community_id: &str) {}
86
87 fn on_community_dissolved(&self, _community_id: &str) {}
90
91 fn buffer_persist(&self, _chat_id: &str, _msg: &Message, _wrapper: Option<([u8; 32], u64)>) -> bool { false }
98}
99
100pub struct NoOpEventHandler;
102impl InboundEventHandler for NoOpEventHandler {}
103
104struct BufferedDm {
108 chat_id: String,
109 msg: Message,
110 wrapper: Option<([u8; 32], u64)>,
111}
112
113pub struct BatchingPersist<'a> {
122 inner: &'a dyn InboundEventHandler,
123 buf: std::sync::Mutex<Vec<BufferedDm>>,
124}
125
126impl<'a> BatchingPersist<'a> {
127 pub fn new(inner: &'a dyn InboundEventHandler) -> Self {
128 Self { inner, buf: std::sync::Mutex::new(Vec::new()) }
129 }
130
131 pub fn buffered(&self) -> usize {
133 self.buf.lock().map(|b| b.len()).unwrap_or(0)
134 }
135
136 pub async fn flush(&self, session: &crate::state::SessionGuard) -> usize {
141 self.try_flush(session).await.unwrap_or(0)
142 }
143
144 pub async fn try_flush(&self, session: &crate::state::SessionGuard) -> Result<usize, String> {
149 let mut drained: Vec<BufferedDm> = match self.buf.lock() {
150 Ok(mut b) => b.drain(..).collect(),
151 Err(_) => return Ok(0),
152 };
153 if drained.is_empty() || !session.is_valid() {
154 return Ok(0);
155 }
156 drained.retain(|e| !crate::state::was_message_deleted(&e.msg.id));
164 if drained.is_empty() {
165 return Ok(0);
166 }
167 let mut groups: Vec<(String, Vec<(&Message, Option<([u8; 32], u64)>)>)> = Vec::new();
169 for e in &drained {
170 match groups.iter_mut().find(|(c, _)| c == &e.chat_id) {
171 Some((_, v)) => v.push((&e.msg, e.wrapper)),
172 None => groups.push((e.chat_id.clone(), vec![(&e.msg, e.wrapper)])),
173 }
174 }
175 match crate::db::events::save_messages_batch_multi(&groups, Some(session)).await {
176 Ok(n) => Ok(n),
177 Err(e) => {
178 crate::log_warn!("[Sync] batched persist failed ({} msgs): {}", drained.len(), e);
179 Err(e)
180 }
181 }
182 }
183}
184
185impl InboundEventHandler for BatchingPersist<'_> {
186 fn buffer_persist(&self, chat_id: &str, msg: &Message, wrapper: Option<([u8; 32], u64)>) -> bool {
187 match self.buf.lock() {
188 Ok(mut b) => {
189 b.push(BufferedDm { chat_id: chat_id.to_string(), msg: msg.clone(), wrapper });
190 true
191 }
192 Err(_) => false,
194 }
195 }
196
197 fn on_dm_received(&self, chat_id: &str, msg: &Message, is_new: bool) {
198 self.inner.on_dm_received(chat_id, msg, is_new)
199 }
200 fn on_file_received(&self, chat_id: &str, msg: &Message, is_new: bool) {
201 self.inner.on_file_received(chat_id, msg, is_new)
202 }
203 fn on_reaction_received(&self, chat_id: &str, msg: &Message) {
204 self.inner.on_reaction_received(chat_id, msg)
205 }
206 fn on_message_deleted(&self, chat_id: &str, message_id: &str) {
207 if let Ok(mut b) = self.buf.lock() {
211 b.retain(|e| e.msg.id != message_id);
212 }
213 self.inner.on_message_deleted(chat_id, message_id)
214 }
215 fn on_community_invite(&self, community_id: &str) {
216 self.inner.on_community_invite(community_id)
217 }
218 fn on_community_message(&self, chat_id: &str, msg: &Message, is_new: bool) {
219 self.inner.on_community_message(chat_id, msg, is_new)
220 }
221 fn on_community_update(&self, chat_id: &str, target_id: &str, msg: &Message) {
222 self.inner.on_community_update(chat_id, target_id, msg)
223 }
224 fn on_community_removed(&self, chat_id: &str, target_id: &str) {
225 self.inner.on_community_removed(chat_id, target_id)
226 }
227 fn on_community_presence(
228 &self,
229 chat_id: &str,
230 npub: &str,
231 joined: bool,
232 event_id: &str,
233 created_at: u64,
234 invited_by: Option<&str>,
235 invited_label: Option<&str>,
236 ) {
237 self.inner.on_community_presence(chat_id, npub, joined, event_id, created_at, invited_by, invited_label)
238 }
239 fn on_community_typing(&self, chat_id: &str, npub: &str, until: u64) {
240 self.inner.on_community_typing(chat_id, npub, until)
241 }
242 fn on_community_webxdc(
243 &self,
244 chat_id: &str,
245 npub: &str,
246 topic_id: &str,
247 node_addr: Option<&str>,
248 event_id: &str,
249 created_at: u64,
250 ) {
251 self.inner.on_community_webxdc(chat_id, npub, topic_id, node_addr, event_id, created_at)
252 }
253 fn on_community_self_removed(&self, community_id: &str) {
254 self.inner.on_community_self_removed(community_id)
255 }
256 fn on_community_refreshed(&self, community_id: &str) {
257 self.inner.on_community_refreshed(community_id)
258 }
259 fn on_community_dissolved(&self, community_id: &str) {
260 self.inner.on_community_dissolved(community_id)
261 }
262}
263
264pub enum PreparedEvent {
266 Processed {
268 result: RumorProcessingResult,
269 contact: String,
270 sender: PublicKey,
271 is_mine: bool,
272 wrapper_event_id: String,
273 wrapper_event_id_bytes: [u8; 32],
274 wrapper_created_at: u64,
275 unwrap_ns: u64,
277 parse_ns: u64,
279 },
280 CommunityInvite {
282 invite: crate::community::invite::CommunityInvite,
283 inviter: String,
285 is_mine: bool,
286 wrapper_event_id_bytes: [u8; 32],
287 wrapper_created_at: u64,
288 rumor_created_at: u64,
292 expires_at: u64,
294 },
295 CommunityInviteV2 {
298 bundle_json: String,
299 community_id: String,
300 inviter: String,
302 is_mine: bool,
303 wrapper_event_id_bytes: [u8; 32],
304 wrapper_created_at: u64,
305 rumor_created_at: u64,
307 expires_at: u64,
309 },
310 DedupSkip {
312 wrapper_id_bytes: [u8; 32],
313 wrapper_created_at: u64,
314 },
315 ErrorSkip {
317 wrapper_id_bytes: [u8; 32],
318 wrapper_created_at: u64,
319 },
320}
321
322pub async fn prepare_event(
327 event: Event,
328 _client: &Client,
329 my_public_key: PublicKey,
330) -> PreparedEvent {
331 let wrapper_created_at = event.created_at.as_secs();
332 let wrapper_event_id_bytes: [u8; 32] = event.id.to_bytes();
333 let wrapper_event_id = event.id.to_hex();
334
335 {
337 let cache = WRAPPER_ID_CACHE.lock().await;
338 if cache.contains(&wrapper_event_id_bytes) {
339 return PreparedEvent::DedupSkip { wrapper_id_bytes: wrapper_event_id_bytes, wrapper_created_at };
340 }
341 }
342
343 if let Ok(true) = crate::db::events::wrapper_event_exists(&wrapper_event_id) {
344 return PreparedEvent::DedupSkip { wrapper_id_bytes: wrapper_event_id_bytes, wrapper_created_at };
345 }
346
347 let unwrap_start = std::time::Instant::now();
349 let signer = match crate::signer::active_signer() {
350 Ok(s) => s,
351 Err(_) => return PreparedEvent::ErrorSkip {
352 wrapper_id_bytes: wrapper_event_id_bytes, wrapper_created_at,
353 },
354 };
355 let (rumor, sender) = match UnwrappedGift::from_gift_wrap_async(&signer, &event).await {
356 Ok(UnwrappedGift { rumor, sender }) => (rumor, sender),
357 Err(_) => return PreparedEvent::ErrorSkip {
358 wrapper_id_bytes: wrapper_event_id_bytes, wrapper_created_at,
359 },
360 };
361
362 let unwrap_ns = unwrap_start.elapsed().as_nanos() as u64;
363
364 let rumor_created_at = rumor.created_at.as_secs();
367
368 let is_mine = sender == my_public_key;
369 let contact = if is_mine {
370 rumor.tags.public_keys().next()
371 .and_then(|pk| pk.to_bech32().ok())
372 .unwrap_or_else(|| sender.to_bech32().unwrap_or_default())
373 } else {
374 sender.to_bech32().unwrap_or_default()
375 };
376
377 if rumor.tags.public_keys().count() > 1 {
379 return PreparedEvent::ErrorSkip {
380 wrapper_id_bytes: wrapper_event_id_bytes, wrapper_created_at,
381 };
382 }
383
384 if rumor.kind == Kind::Custom(crate::stored_event::event_kind::COMMUNITY_INVITE_BUNDLE) {
387 return match crate::community::invite::parse_invite_rumor(rumor.kind, &rumor.content) {
388 Some(invite) => PreparedEvent::CommunityInvite {
389 invite, inviter: contact.clone(), is_mine, wrapper_event_id_bytes, wrapper_created_at, rumor_created_at,
390 expires_at: crate::community::invite::expiration_secs(&rumor.tags).unwrap_or(0),
391 },
392 None => PreparedEvent::ErrorSkip { wrapper_id_bytes: wrapper_event_id_bytes, wrapper_created_at },
393 };
394 }
395
396 if rumor.kind == Kind::Custom(crate::community::v2::kind::DIRECT_INVITE) {
399 return match crate::community::v2::invite::CommunityInvite::from_bundle_json(&rumor.content)
400 .ok()
401 .and_then(|b| serde_json::to_string(&b).ok().map(|j| (b.community_id, j)))
402 {
403 Some((community_id, bundle_json)) => PreparedEvent::CommunityInviteV2 {
404 community_id,
405 bundle_json,
406 inviter: contact.clone(), is_mine,
408 wrapper_event_id_bytes,
409 wrapper_created_at,
410 rumor_created_at,
411 expires_at: crate::community::invite::expiration_secs(&rumor.tags).unwrap_or(0),
412 },
413 None => PreparedEvent::ErrorSkip { wrapper_id_bytes: wrapper_event_id_bytes, wrapper_created_at },
414 };
415 }
416
417 let Some(rumor_id) = rumor.id else {
419 return PreparedEvent::ErrorSkip {
420 wrapper_id_bytes: wrapper_event_id_bytes, wrapper_created_at,
421 };
422 };
423
424 let rumor_event = RumorEvent {
425 id: rumor_id,
426 kind: rumor.kind,
427 content: rumor.content,
428 tags: rumor.tags,
429 created_at: rumor.created_at,
430 pubkey: rumor.pubkey,
431 };
432 let rumor_context = RumorContext {
433 sender,
434 is_mine,
435 conversation_id: contact.clone(),
436 conversation_type: ConversationType::DirectMessage,
437 };
438
439 let parse_start = std::time::Instant::now();
440 let download_dir = crate::db::get_download_dir();
441 match process_rumor(rumor_event, rumor_context, &download_dir) {
442 Ok(result) => {
443 let parse_ns = parse_start.elapsed().as_nanos() as u64;
444 PreparedEvent::Processed {
445 result, contact, sender, is_mine,
446 wrapper_event_id, wrapper_event_id_bytes, wrapper_created_at,
447 unwrap_ns, parse_ns,
448 }
449 }
450 Err(e) => {
451 log_warn!("[EventHandler] Failed to process rumor: {}", e);
452 PreparedEvent::ErrorSkip {
453 wrapper_id_bytes: wrapper_event_id_bytes, wrapper_created_at,
454 }
455 }
456 }
457}
458
459pub const DIRECT_INVITE_LIFETIME_SECS: u64 = 24 * 3600 + 3600;
480
481fn expired_invite(expires_at: u64, rumor_created_at: u64) -> bool {
482 let now = nostr_sdk::prelude::Timestamp::now().as_secs();
483 if expires_at != 0 && expires_at <= now {
484 return true;
485 }
486 rumor_created_at.saturating_add(DIRECT_INVITE_LIFETIME_SECS) <= now
487}
488
489#[cfg(test)]
490mod invite_expiry_tests {
491 use super::*;
492
493 fn now() -> u64 {
494 nostr_sdk::prelude::Timestamp::now().as_secs()
495 }
496
497 #[test]
498 fn declared_deadline_is_honored() {
499 assert!(expired_invite(now() - 10, now()));
500 assert!(!expired_invite(now() + 3600, now()));
501 }
502
503 #[test]
504 fn tagless_invites_die_at_the_recipient_lifetime() {
505 assert!(!expired_invite(0, now() - 3600));
507 assert!(expired_invite(0, now() - DIRECT_INVITE_LIFETIME_SECS));
510 assert!(expired_invite(0, now() - 90 * 24 * 3600));
511 }
512
513 #[test]
514 fn lifetime_caps_a_generous_declared_deadline() {
515 assert!(expired_invite(now() + 7 * 24 * 3600, now() - DIRECT_INVITE_LIFETIME_SECS));
517 }
518}
519
520pub async fn commit_prepared_event(
521 prepared: PreparedEvent,
522 is_new: bool,
523 handler: &dyn InboundEventHandler,
524) -> bool {
525 let session = crate::state::SessionGuard::capture();
526 if !session.is_valid() {
527 return false;
528 }
529 match prepared {
530 PreparedEvent::Processed { result, contact, sender, is_mine, wrapper_event_id, wrapper_event_id_bytes, wrapper_created_at, .. } => {
531 {
533 let mut cache = WRAPPER_ID_CACHE.lock().await;
534 cache.insert(wrapper_event_id_bytes);
535 }
536
537 if !is_mine {
540 let blocked = {
541 let state = crate::state::STATE.lock().await;
542 state.get_profile(&contact).map_or(false, |p| p.flags.is_blocked())
543 };
544 if blocked {
545 let _ = crate::db::wrappers::save_processed_wrapper(&wrapper_event_id_bytes, wrapper_created_at, crate::db::wrappers::TRANSPORT_NIP17);
546 return false;
547 }
548 }
549
550 if !matches!(result, RumorProcessingResult::TextMessage(_) | RumorProcessingResult::FileAttachment(_)) {
556 let _ = crate::db::wrappers::save_processed_wrapper(&wrapper_event_id_bytes, wrapper_created_at, crate::db::wrappers::TRANSPORT_NIP17);
557 }
558
559 match result {
560 RumorProcessingResult::TextMessage(mut msg) => {
561 msg.wrapper_event_id = Some(wrapper_event_id.clone());
562 commit_dm_message(msg, &contact, is_mine, is_new, &wrapper_event_id, wrapper_event_id_bytes, wrapper_created_at, handler, false).await
563 }
564 RumorProcessingResult::FileAttachment(mut msg) => {
565 msg.wrapper_event_id = Some(wrapper_event_id.clone());
566 if !is_mine {
581 for att in &mut msg.attachments {
582 if att.size == 0
583 && (att.url.starts_with("https://") || att.url.starts_with("http://"))
584 {
585 if let Ok(Some(size)) = tokio::time::timeout(
586 std::time::Duration::from_secs(3),
587 crate::net::get_remote_file_size(&att.url),
588 ).await {
589 att.size = size;
590 }
591 }
592 }
593 }
594 commit_dm_message(msg, &contact, is_mine, is_new, &wrapper_event_id, wrapper_event_id_bytes, wrapper_created_at, handler, true).await
595 }
596 RumorProcessingResult::Reaction(reaction) => {
597 commit_reaction(reaction, &contact, is_mine, &wrapper_event_id, handler).await
598 }
599 RumorProcessingResult::Edit { message_id, new_content, edited_at, emoji_tags, mut event } => {
600 commit_edit(&mut event, &contact, &message_id, &new_content, edited_at, emoji_tags, &wrapper_event_id).await
601 }
602 RumorProcessingResult::TypingIndicator { profile_id, until } => {
603 let active_typers = {
604 let mut state = crate::state::STATE.lock().await;
605 state.update_typing_and_get_active(&contact, &profile_id, until)
606 };
607 crate::traits::emit_event("typing-update", &serde_json::json!({
608 "conversation_id": contact,
609 "typers": active_typers,
610 }));
611 false
612 }
613 RumorProcessingResult::PivxPayment { gift_code, amount_piv, address, message_id, mut event } => {
614 if crate::db::events::event_exists(&event.id).unwrap_or(false) {
615 return false;
616 }
617 event.wrapper_event_id = Some(wrapper_event_id.clone());
618 let ts = event.created_at;
619 let _ = crate::db::events::save_pivx_payment_event(&contact, event).await;
620 crate::traits::emit_event("pivx_payment_received", &serde_json::json!({
621 "conversation_id": contact,
622 "gift_code": gift_code, "amount_piv": amount_piv,
623 "address": address, "message_id": message_id,
624 "sender": sender.to_hex(), "is_mine": is_mine,
625 "at": ts * 1000,
626 }));
627 true
628 }
629 RumorProcessingResult::UnknownEvent(mut event) => {
630 event.wrapper_event_id = Some(wrapper_event_id.clone());
631 if let Ok(chat_id) = crate::db::id_cache::get_or_create_chat_id(&contact) {
633 event.chat_id = chat_id;
634 }
635 let _ = crate::db::events::save_event(&event).await;
636 false
637 }
638 RumorProcessingResult::LeaveRequest { .. } => false,
639 RumorProcessingResult::WebxdcPeerAdvertisement { .. } |
640 RumorProcessingResult::WebxdcPeerLeft { .. } => {
641 false
643 }
644 RumorProcessingResult::WallpaperChanged {
645 sender_npub, created_at, url, decryption_key, decryption_nonce,
646 plaintext_hash, mime, blur, dim, event_id,
647 } => {
648 let _ = crate::wallpaper::apply_received_wallpaper(
649 &contact, &sender_npub, created_at, &url,
650 &decryption_key, &decryption_nonce,
651 plaintext_hash.as_deref(), mime.as_deref(),
652 blur, dim,
653 &event_id,
654 ).await;
655 true
658 }
659 RumorProcessingResult::DeletionRequest { target_event_id } => {
660 if commit_deletion(&target_event_id, &contact, &sender, handler).await {
664 true
665 } else {
666 commit_reaction_deletion(&target_event_id, &sender).await
667 }
668 }
669 RumorProcessingResult::Ignored => false,
670 }
671 }
672 PreparedEvent::CommunityInvite { invite, inviter, is_mine, wrapper_event_id_bytes, wrapper_created_at, rumor_created_at, expires_at } => {
673 {
676 let mut cache = WRAPPER_ID_CACHE.lock().await;
677 cache.insert(wrapper_event_id_bytes);
678 }
679 let _ = crate::db::wrappers::save_processed_wrapper(&wrapper_event_id_bytes, wrapper_created_at, crate::db::wrappers::TRANSPORT_NIP17);
680
681 if is_mine {
683 return false;
684 }
685
686 if expired_invite(expires_at, rumor_created_at) {
690 return false;
691 }
692
693 if let Err(e) = invite.validate() {
696 log_warn!("[community] invite rejected: {}", e);
697 return false;
698 }
699
700 let community_id = invite.community_id.clone();
705 let already_held = crate::community::CommunityId(
708 match crate::simd::hex::hex_to_bytes_32_checked(&community_id) {
709 Some(b) => b,
710 None => { log_warn!("[community] invite has malformed id"); return false; }
711 },
712 );
713 if crate::db::community::community_exists(&already_held).unwrap_or(false) {
714 return false;
715 }
716 if crate::db::community::pending_invite_exists(&community_id).unwrap_or(false) {
717 return false;
718 }
719
720 if crate::community::list::tombstone_suppresses(&community_id, rumor_created_at) {
726 return false;
727 }
728
729 let bundle_json = match invite.to_json() {
732 Ok(j) => j,
733 Err(e) => { log_warn!("[community] invite re-serialize failed: {}", e); return false; }
734 };
735 match crate::db::community::save_pending_invite(&community_id, &bundle_json, &inviter, expires_at as i64) {
736 Ok(true) => {
737 handler.on_community_invite(&community_id);
738 let invite_warm = invite.clone();
742 let bg = crate::state::SessionGuard::capture();
743 tokio::spawn(async move {
744 if !bg.is_valid() {
745 return;
746 }
747 crate::community::service::preload_community(&invite_warm).await;
748 });
749 }
750 Ok(false) => {} Err(e) => log_warn!("[community] invite park failed: {}", e),
752 }
753 false
754 }
755 PreparedEvent::CommunityInviteV2 { bundle_json, community_id, inviter, is_mine, wrapper_event_id_bytes, wrapper_created_at, rumor_created_at, expires_at } => {
756 {
757 let mut cache = WRAPPER_ID_CACHE.lock().await;
758 cache.insert(wrapper_event_id_bytes);
759 }
760 let _ = crate::db::wrappers::save_processed_wrapper(&wrapper_event_id_bytes, wrapper_created_at, crate::db::wrappers::TRANSPORT_NIP17);
761
762 if is_mine {
764 return false;
765 }
766 if expired_invite(expires_at, rumor_created_at) {
768 return false;
769 }
770 let held = match crate::simd::hex::hex_to_bytes_32_checked(&community_id) {
773 Some(b) => crate::community::CommunityId(b),
774 None => return false,
775 };
776 if crate::db::community::community_exists(&held).unwrap_or(false) {
777 return false;
778 }
779 if crate::db::community::pending_invite_exists(&community_id).unwrap_or(false) {
780 return false;
781 }
782 if crate::community::list::tombstone_suppresses(&community_id, rumor_created_at) {
787 return false;
788 }
789 match crate::db::community::save_pending_invite(&community_id, &bundle_json, &inviter, expires_at as i64) {
791 Ok(true) => handler.on_community_invite(&community_id),
792 Ok(false) => {} Err(e) => log_warn!("[community] v2 invite park failed: {}", e),
794 }
795 false
796 }
797 PreparedEvent::DedupSkip { wrapper_id_bytes, wrapper_created_at } => {
798 if wrapper_created_at > 0 {
805 if crate::db::wrappers::processed_wrapper_exists(&wrapper_id_bytes) {
806 let _ = crate::db::wrappers::update_wrapper_timestamp(&wrapper_id_bytes, wrapper_created_at);
807 } else {
808 let wrapper_hex = crate::simd::hex::bytes_to_hex_32(&wrapper_id_bytes);
809 if crate::db::events::wrapper_event_exists(&wrapper_hex).unwrap_or(false) {
810 let _ = crate::db::wrappers::save_processed_wrapper(&wrapper_id_bytes, wrapper_created_at, crate::db::wrappers::TRANSPORT_NIP17);
811 }
812 }
813 }
814 false
815 }
816 PreparedEvent::ErrorSkip { wrapper_id_bytes, wrapper_created_at } => {
817 let _ = crate::db::wrappers::save_processed_wrapper(&wrapper_id_bytes, wrapper_created_at, crate::db::wrappers::TRANSPORT_NIP17);
818 false
819 }
820 }
821}
822
823#[allow(clippy::too_many_arguments)]
831async fn commit_dm_message(
832 mut msg: Message,
833 contact: &str,
834 _is_mine: bool,
835 is_new: bool,
836 wrapper_event_id: &str,
837 wrapper_event_id_bytes: [u8; 32],
838 wrapper_created_at: u64,
839 handler: &dyn InboundEventHandler,
840 is_file: bool,
841) -> bool {
842 let ledger_wrapper = || {
843 let _ = crate::db::wrappers::save_processed_wrapper(
844 &wrapper_event_id_bytes, wrapper_created_at, crate::db::wrappers::TRANSPORT_NIP17,
845 );
846 };
847 if let Ok(true) = crate::db::events::message_exists_in_db(&msg.id) {
849 if let Ok(updated) = crate::db::events::update_wrapper_event_id(&msg.id, wrapper_event_id) {
851 if !updated {
852 let mut cache = WRAPPER_ID_CACHE.lock().await;
853 cache.insert(wrapper_event_id_bytes);
854 }
855 }
856 ledger_wrapper();
857 return false;
858 }
859
860 if !msg.replied_to.is_empty() {
862 let _ = crate::db::events::populate_reply_context(&mut msg).await;
863 }
864
865 let added = {
867 let mut state = crate::state::STATE.lock().await;
868 let added = state.add_message_to_participant(contact, &msg);
869 if is_file && added {
870 state.update_typing_and_get_active(contact, contact, 0);
871 }
872 added
873 };
874
875 if added {
876 crate::traits::emit_event("message_new", &serde_json::json!({
878 "message": &msg,
879 "chat_id": contact
880 }));
881
882 if is_file {
884 handler.on_file_received(contact, &msg, is_new);
885 } else {
886 handler.on_dm_received(contact, &msg, is_new);
887 }
888
889 if !handler.buffer_persist(contact, &msg, Some((wrapper_event_id_bytes, wrapper_created_at))) {
894 if crate::db::events::save_message(contact, &msg).await.is_ok() {
895 ledger_wrapper();
896 }
897 }
898 } else {
899 ledger_wrapper();
902 }
903
904 added
905}
906
907async fn commit_reaction(
909 reaction: crate::types::Reaction,
910 contact: &str,
911 is_mine: bool,
912 wrapper_event_id: &str,
913 handler: &dyn InboundEventHandler,
914) -> bool {
915 let msg_for_emit = {
917 let mut state = crate::state::STATE.lock().await;
918 if let Some((chat_id, was_added)) = state.add_reaction_to_message(&reaction.reference_id, reaction.clone()) {
919 if was_added {
920 state.find_message(&reaction.reference_id)
921 .map(|(_, msg)| (chat_id, msg))
922 } else { None }
923 } else { None }
924 };
925
926 if let Some((chat_id, mut msg)) = msg_for_emit {
927 crate::traits::emit_message_update(&chat_id, &reaction.reference_id, &mut msg).await;
928 let _ = crate::db::events::save_message(&chat_id, &msg).await;
929 handler.on_reaction_received(&chat_id, &msg);
930 }
931
932 if let Ok(chat_id) = crate::db::id_cache::get_chat_id_by_identifier(contact) {
934 let _ = crate::db::events::save_reaction_event(
935 &reaction, chat_id, None, is_mine, Some(wrapper_event_id.to_string())
936 ).await;
937 }
938
939 true
940}
941
942async fn commit_edit(
944 event: &mut crate::stored_event::StoredEvent,
945 contact: &str,
946 message_id: &str,
947 new_content: &str,
948 edited_at: u64,
949 emoji_tags: Vec<crate::types::EmojiTag>,
950 wrapper_event_id: &str,
951) -> bool {
952 if crate::db::events::event_exists(&event.id).unwrap_or(false) {
953 return false;
954 }
955 if let Ok(chat_id) = crate::db::id_cache::get_chat_id_by_identifier(contact) {
956 event.chat_id = chat_id;
957 }
958 event.wrapper_event_id = Some(wrapper_event_id.to_string());
959 let _ = crate::db::events::save_event(event).await;
960
961 let msg_for_emit = {
962 let mut state = crate::state::STATE.lock().await;
963 state.update_message_in_chat(contact, message_id, |msg| {
964 msg.apply_edit(new_content.to_string(), edited_at, emoji_tags.clone());
965 })
966 };
967 if let Some(mut msg) = msg_for_emit {
968 crate::traits::emit_message_update(contact, message_id, &mut msg).await;
969 }
970 true
971}
972
973async fn commit_deletion(
986 target_event_id: &str,
987 contact: &str,
988 sender: &PublicKey,
989 handler: &dyn InboundEventHandler,
990) -> bool {
991 let (mine, chat_id) = {
1009 let state = crate::state::STATE.lock().await;
1010 match state.find_message(target_event_id) {
1011 Some((chat, msg)) => (msg.mine, chat.id.clone()),
1012 None => return false,
1013 }
1014 };
1015
1016 let authorized = if mine {
1027 match crate::state::my_public_key() {
1028 Some(my_pk) => *sender == my_pk,
1029 None => false,
1030 }
1031 } else {
1032 match nostr_sdk::prelude::PublicKey::from_bech32(&chat_id) {
1033 Ok(counterpart) => sender == &counterpart,
1034 Err(_) => false, }
1036 };
1037 if !authorized {
1038 eprintln!(
1039 "[NIP-17 cooperative-delete] unauthorized: sender {} not the author of target {} (mine={}, chat={})",
1040 sender.to_hex(), target_event_id, mine, chat_id
1041 );
1042 return false;
1043 }
1044
1045 let removed = {
1047 let mut state = crate::state::STATE.lock().await;
1048 state.remove_message(target_event_id)
1049 };
1050 let removed_msg = match removed {
1051 Some((_chat_id, msg)) => msg,
1052 None => return false,
1053 };
1054 crate::state::note_message_deleted(target_event_id);
1057
1058 let unique = crate::deletion::filter_unreferenced_attachments(
1068 target_event_id,
1069 removed_msg.attachments,
1070 ).await;
1071 crate::deletion::delete_cached_attachment_files_pub(&unique);
1072
1073 if let Err(e) = crate::db::events::delete_event(target_event_id).await {
1075 eprintln!(
1076 "[NIP-17 cooperative-delete] DB delete failed for {}: {}",
1077 target_event_id, e
1078 );
1079 }
1080
1081 crate::traits::emit_event(
1084 "message_removed",
1085 &serde_json::json!({
1086 "id": target_event_id,
1087 "chat_id": &chat_id,
1088 "reason": "deleted-by-sender",
1089 }),
1090 );
1091
1092 handler.on_message_deleted(&chat_id, target_event_id);
1093 let _ = contact;
1094 true
1095}
1096
1097async fn commit_reaction_deletion(target_reaction_id: &str, sender: &PublicKey) -> bool {
1102 let found = {
1103 let state = crate::state::STATE.lock().await;
1104 state.find_reaction(target_reaction_id)
1105 };
1106 let (chat_id, message_id, author_npub, _is_community) = match found {
1107 Some(v) => v,
1108 None => return false,
1109 };
1110
1111 let authorized = nostr_sdk::prelude::PublicKey::parse(&author_npub)
1113 .map(|pk| pk == *sender)
1114 .unwrap_or(false);
1115 if !authorized {
1116 eprintln!(
1117 "[reaction-delete] unauthorized: sender {} is not the author of reaction {}",
1118 sender.to_hex(), target_reaction_id
1119 );
1120 return false;
1121 }
1122
1123 let updated = {
1124 let mut state = crate::state::STATE.lock().await;
1125 state.remove_reaction_from_message(&message_id, target_reaction_id)
1126 };
1127 let mut message = match updated {
1128 Some((_cid, msg)) => msg,
1129 None => return false,
1130 };
1131
1132 if let Err(e) = crate::db::events::delete_event(target_reaction_id).await {
1135 eprintln!("[reaction-delete] DB delete failed for {}: {}", target_reaction_id, e);
1136 }
1137
1138 crate::traits::emit_message_update(&chat_id, &message_id, &mut message).await;
1139 true
1140}
1141
1142pub async fn process_event(
1151 event: Event,
1152 is_new: bool,
1153 handler: &dyn InboundEventHandler,
1154) -> std::result::Result<bool, String> {
1155 let client = crate::state::nostr_client()
1156 .ok_or_else(|| "Nostr client not initialized".to_string())?;
1157 let my_pk = crate::state::my_public_key()
1158 .ok_or_else(|| "Public key not initialized".to_string())?;
1159 let prepared = prepare_event(event, &client, my_pk).await;
1160 Ok(commit_prepared_event(prepared, is_new, handler).await)
1161}