1use crate::stored_event::{StoredEvent, event_kind};
4use rusqlite::OptionalExtension;
5use crate::crypto::maybe_encrypt;
6use crate::types::{Message, Attachment, Reaction};
7
8async fn encrypt_event_content(event: &StoredEvent) -> String {
16 if event.kind == event_kind::CHAT_MESSAGE
17 || event.kind == event_kind::PRIVATE_DIRECT_MESSAGE
18 || event.kind == event_kind::MESSAGE_EDIT
19 {
20 maybe_encrypt(event.content.clone()).await
21 } else {
22 event.content.clone()
23 }
24}
25
26fn insert_event_row(conn: &rusqlite::Connection, event: &StoredEvent, content: &str, tags_json: &str) -> Result<(), String> {
33 let mut stmt = conn.prepare_cached(
36 r#"
37 INSERT INTO events (
38 id, kind, chat_id, user_id, content, tags, reference_id,
39 created_at, received_at, mine, pending, failed, wrapper_event_id, npub, preview_metadata
40 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15)
41 ON CONFLICT(id) DO UPDATE SET
42 kind = excluded.kind, chat_id = excluded.chat_id, user_id = excluded.user_id,
43 content = excluded.content, tags = excluded.tags, reference_id = excluded.reference_id,
44 created_at = excluded.created_at, received_at = excluded.received_at,
45 mine = excluded.mine, pending = excluded.pending, failed = excluded.failed,
46 wrapper_event_id = COALESCE(excluded.wrapper_event_id, events.wrapper_event_id),
47 npub = excluded.npub, preview_metadata = excluded.preview_metadata
48 "#,
49 ).map_err(|e| format!("prepare save event: {}", e))?;
50 stmt.execute(
51 rusqlite::params![
52 event.id, event.kind as i32, event.chat_id, event.user_id, content, tags_json,
53 event.reference_id, event.created_at as i64, event.received_at as i64,
54 event.mine as i32, event.pending as i32, event.failed as i32,
55 event.wrapper_event_id, event.npub, event.preview_metadata,
56 ],
57 ).map_err(|e| format!("Failed to save event: {}", e))?;
58 Ok(())
59}
60
61pub async fn save_event(event: &StoredEvent) -> Result<(), String> {
62 let tags_json = serde_json::to_string(&event.tags).unwrap_or_else(|_| "[]".to_string());
63 let content = encrypt_event_content(event).await;
64 let conn = super::get_write_connection_guard_static()?;
65 insert_event_row(&conn, event, &content, &tags_json)
66}
67
68fn extract_bot_tags(tags: &[Vec<String>]) -> Vec<String> {
71 tags.iter()
72 .filter(|t| t.len() >= 2 && t[0] == "bot")
73 .map(|t| t[1].clone())
74 .collect()
75}
76
77fn extract_expiration_tag(tags: &[Vec<String>]) -> Option<u64> {
80 tags.iter()
81 .find(|t| t.len() >= 2 && t[0] == "expiration")
82 .and_then(|t| t[1].parse::<u64>().ok())
83}
84
85pub fn event_exists(event_id: &str) -> Result<bool, String> {
87 let conn = super::get_db_connection_guard_static()?;
88 event_exists_on(&conn, event_id)
89}
90
91fn event_exists_on(conn: &rusqlite::Connection, event_id: &str) -> Result<bool, String> {
94 let mut stmt = conn.prepare_cached("SELECT EXISTS(SELECT 1 FROM events WHERE id = ?1)")
95 .map_err(|e| format!("prepare event existence: {}", e))?;
96 stmt.query_row(rusqlite::params![event_id], |row| row.get(0))
97 .map_err(|e| format!("Failed to check event existence: {}", e))
98}
99
100fn reaction_to_stored_event(
102 reaction: &Reaction,
103 chat_id: i64,
104 user_id: Option<i64>,
105 mine: bool,
106 wrapper_event_id: Option<String>,
107) -> StoredEvent {
108 let mut tags: Vec<Vec<String>> = vec![
112 vec!["e".to_string(), reaction.reference_id.clone()],
113 ];
114 if let Some(url) = &reaction.emoji_url {
115 if reaction.emoji.starts_with(':') && reaction.emoji.ends_with(':') && reaction.emoji.len() >= 3 {
116 let shortcode = &reaction.emoji[1..reaction.emoji.len() - 1];
117 if !shortcode.is_empty() && !url.is_empty() {
118 tags.push(vec!["emoji".to_string(), shortcode.to_string(), url.clone()]);
119 }
120 }
121 }
122 StoredEvent {
123 id: reaction.id.clone(),
124 kind: event_kind::REACTION,
125 chat_id,
126 user_id,
127 content: reaction.emoji.clone(),
128 tags,
129 reference_id: Some(reaction.reference_id.clone()),
130 created_at: std::time::SystemTime::now()
131 .duration_since(std::time::UNIX_EPOCH)
132 .map(|d| d.as_secs()).unwrap_or(0),
133 received_at: std::time::SystemTime::now()
134 .duration_since(std::time::UNIX_EPOCH)
135 .map(|d| d.as_millis() as u64).unwrap_or(0),
136 mine,
137 pending: false,
138 failed: false,
139 wrapper_event_id,
140 npub: Some(reaction.author_id.clone()),
141 preview_metadata: None,
142 }
143}
144
145pub async fn save_reaction_event(
147 reaction: &Reaction,
148 chat_id: i64,
149 user_id: Option<i64>,
150 mine: bool,
151 wrapper_event_id: Option<String>,
152) -> Result<(), String> {
153 let event = reaction_to_stored_event(reaction, chat_id, user_id, mine, wrapper_event_id);
154 save_event(&event).await
155}
156
157pub async fn save_message(chat_id: &str, message: &Message) -> Result<(), String> {
166 let chat_int_id = super::id_cache::get_or_create_chat_id(chat_id)?;
167
168 let user_int_id = if let Some(ref npub_str) = message.npub {
169 super::id_cache::get_or_create_user_id(npub_str)?
170 } else {
171 None
172 };
173
174 let event = message_to_stored_event(message, chat_int_id, user_int_id);
175
176 let tags_json = serde_json::to_string(&event.tags).unwrap_or_else(|_| "[]".to_string());
181 let content = encrypt_event_content(&event).await;
182 {
183 let conn = super::get_write_connection_guard_static()?;
184 let tx = conn.unchecked_transaction().map_err(|e| format!("save_message tx: {e}"))?;
185 insert_event_row(&tx, &event, &content, &tags_json)?;
186 super::attachments::insert_attachment_rows(&tx, &message.id, &message.attachments)?;
187 tx.commit().map_err(|e| format!("save_message commit: {e}"))?;
188 }
189
190 for reaction in &message.reactions {
192 if !event_exists(&reaction.id)? {
193 let user_id = super::id_cache::get_or_create_user_id(&reaction.author_id)?;
194 let is_mine = super::get_current_account()
195 .map(|npub| reaction.author_id == npub)
196 .unwrap_or(false);
197 save_reaction_event(reaction, chat_int_id, user_id, is_mine, None).await?;
198 }
199 }
200
201 Ok(())
202}
203
204struct BatchRow<'a> {
208 message: &'a Message,
209 event: StoredEvent,
210 content: String,
211 tags_json: String,
212 reactions: Vec<(StoredEvent, String)>,
213 wrapper: Option<([u8; 32], u64)>,
217}
218
219async fn prepare_batch_rows<'a>(
223 chat_id: &str,
224 messages: &[(&'a Message, Option<([u8; 32], u64)>)],
225 rows: &mut Vec<BatchRow<'a>>,
226) -> Result<(), String> {
227 let chat_int_id = super::id_cache::get_or_create_chat_id(chat_id)?;
228 let my_npub = super::get_current_account();
229 for (message, wrapper) in messages {
230 let user_int_id = match &message.npub {
231 Some(npub_str) => super::id_cache::get_or_create_user_id(npub_str)?,
232 None => None,
233 };
234 let event = message_to_stored_event(message, chat_int_id, user_int_id);
235 let tags_json = serde_json::to_string(&event.tags).unwrap_or_else(|_| "[]".to_string());
236 let content = encrypt_event_content(&event).await;
237 let mut reactions: Vec<(StoredEvent, String)> = Vec::with_capacity(message.reactions.len());
238 for reaction in &message.reactions {
239 let user_id = super::id_cache::get_or_create_user_id(&reaction.author_id)?;
240 let is_mine = my_npub.as_deref().map(|n| reaction.author_id == n).unwrap_or(false);
241 let rev = reaction_to_stored_event(reaction, chat_int_id, user_id, is_mine, None);
242 let rtags = serde_json::to_string(&rev.tags).unwrap_or_else(|_| "[]".to_string());
243 reactions.push((rev, rtags));
244 }
245 rows.push(BatchRow { message, event, content, tags_json, reactions, wrapper: *wrapper });
246 }
247 Ok(())
248}
249
250fn write_batch_rows(rows: &[BatchRow<'_>]) -> Result<usize, String> {
256 let conn = super::get_write_connection_guard_static()?;
257 let tx = conn.unchecked_transaction().map_err(|e| format!("batch tx: {e}"))?;
258 let mut saved = 0usize;
259 for row in rows {
260 tx.execute_batch("SAVEPOINT batch_row").map_err(|e| format!("batch savepoint: {e}"))?;
265 let row_written = insert_event_row(&tx, &row.event, &row.content, &row.tags_json)
266 .and_then(|_| super::attachments::insert_attachment_rows(&tx, &row.message.id, &row.message.attachments));
267 if let Err(e) = row_written {
268 crate::log_warn!("[DB] batch skip {}: {}", &row.message.id[..8.min(row.message.id.len())], e);
269 let _ = tx.execute_batch("ROLLBACK TO batch_row; RELEASE batch_row");
270 continue;
271 }
272 saved += 1;
273 for (rev, rtags) in &row.reactions {
274 if event_exists_on(&tx, &rev.id).unwrap_or(true) {
277 continue;
278 }
279 if let Err(e) = insert_event_row(&tx, rev, &rev.content, rtags) {
280 crate::log_warn!("[DB] batch reaction {}: {}", &rev.id[..8.min(rev.id.len())], e);
281 }
282 }
283 if let Some((wrapper_id, wrapper_created_at)) = &row.wrapper {
285 let mut stmt = tx.prepare_cached(
286 "INSERT OR IGNORE INTO processed_wrappers (wrapper_id, wrapper_created_at, transport) VALUES (?1, ?2, ?3)",
287 ).map_err(|e| format!("prepare wrapper ledger: {e}"))?;
288 if let Err(e) = stmt.execute(rusqlite::params![
289 &wrapper_id[..], *wrapper_created_at as i64, super::wrappers::TRANSPORT_NIP17,
290 ]) {
291 crate::log_warn!("[DB] batch wrapper ledger {}: {}", &row.message.id[..8.min(row.message.id.len())], e);
292 }
293 }
294 tx.execute_batch("RELEASE batch_row").map_err(|e| format!("batch release: {e}"))?;
295 }
296 tx.commit().map_err(|e| format!("batch commit: {e}"))?;
297 Ok(saved)
298}
299
300pub async fn save_messages_batch(
313 chat_id: &str,
314 messages: &[&Message],
315 session: Option<&crate::state::SessionGuard>,
316) -> Result<usize, String> {
317 if messages.is_empty() {
318 return Ok(0);
319 }
320 let with_wrappers: Vec<(&Message, Option<([u8; 32], u64)>)> =
321 messages.iter().map(|m| (*m, None)).collect();
322 let mut rows = Vec::with_capacity(messages.len());
323 prepare_batch_rows(chat_id, &with_wrappers, &mut rows).await?;
324 if session.is_some_and(|s| !s.is_valid()) {
325 return Ok(0);
326 }
327 write_batch_rows(&rows)
328}
329
330pub async fn save_messages_batch_multi(
335 groups: &[(String, Vec<(&Message, Option<([u8; 32], u64)>)>)],
336 session: Option<&crate::state::SessionGuard>,
337) -> Result<usize, String> {
338 let total: usize = groups.iter().map(|(_, m)| m.len()).sum();
339 if total == 0 {
340 return Ok(0);
341 }
342 let mut rows = Vec::with_capacity(total);
343 for (chat_id, messages) in groups {
344 prepare_batch_rows(chat_id, messages, &mut rows).await?;
345 }
346 if session.is_some_and(|s| !s.is_valid()) {
347 return Ok(0);
348 }
349 write_batch_rows(&rows)
350}
351
352fn message_to_stored_event(message: &Message, chat_id: i64, user_id: Option<i64>) -> StoredEvent {
354 let kind = if !message.attachments.is_empty() {
355 event_kind::FILE_ATTACHMENT
356 } else {
357 event_kind::PRIVATE_DIRECT_MESSAGE
358 };
359
360 let mut tags: Vec<Vec<String>> = Vec::new();
361
362 let ms = message.at % 1000;
364 if ms > 0 {
365 tags.push(vec!["ms".to_string(), ms.to_string()]);
366 }
367
368 if !message.replied_to.is_empty() {
370 tags.push(vec![
371 "e".to_string(),
372 message.replied_to.clone(),
373 "".to_string(),
374 "reply".to_string(),
375 ]);
376 }
377
378 for et in &message.emoji_tags {
383 tags.push(vec!["emoji".to_string(), et.shortcode.clone(), et.url.clone()]);
384 }
385
386 for npub in &message.addressed_bots {
389 tags.push(vec!["bot".to_string(), npub.clone()]);
390 }
391
392 if let Some(exp) = message.expiration {
395 tags.push(vec!["expiration".to_string(), exp.to_string()]);
396 }
397
398 let preview_metadata = message.preview_metadata.as_ref()
399 .and_then(|m| serde_json::to_string(m).ok());
400
401 StoredEvent {
402 id: message.id.clone(),
403 kind,
404 chat_id,
405 user_id,
406 content: message.content.clone(),
407 tags,
408 reference_id: None,
409 created_at: message.at / 1000,
410 received_at: std::time::SystemTime::now()
411 .duration_since(std::time::UNIX_EPOCH)
412 .map(|d| d.as_millis() as u64)
413 .unwrap_or(0),
414 mine: message.mine,
415 pending: message.pending,
416 failed: message.failed,
417 wrapper_event_id: message.wrapper_event_id.clone(),
418 npub: message.npub.clone(),
419 preview_metadata,
420 }
421}
422
423pub async fn save_pivx_payment_event(
425 conversation_id: &str,
426 mut event: StoredEvent,
427) -> Result<(), String> {
428 event.chat_id = super::id_cache::get_or_create_chat_id(conversation_id)?;
429 save_event(&event).await
430}
431
432pub async fn save_system_event_by_id(
435 event_id: &str,
436 conversation_id: &str,
437 event_type: crate::stored_event::SystemEventType,
438 member_npub: &str,
439 member_name: Option<&str>,
440) -> Result<bool, String> {
441 let now_secs = std::time::SystemTime::now()
442 .duration_since(std::time::UNIX_EPOCH)
443 .map(|d| d.as_secs()).unwrap_or(0);
444 save_system_event_at(event_id, conversation_id, event_type, member_npub, member_name, now_secs, None, None).await
445}
446
447pub async fn save_system_event_at(
451 event_id: &str,
452 conversation_id: &str,
453 event_type: crate::stored_event::SystemEventType,
454 member_npub: &str,
455 member_name: Option<&str>,
456 created_at_secs: u64,
457 invited_by: Option<&str>,
460 invited_label: Option<&str>,
461) -> Result<bool, String> {
462 let chat_id = super::id_cache::get_or_create_chat_id(conversation_id)?;
463
464 let now_secs = std::time::SystemTime::now()
465 .duration_since(std::time::UNIX_EPOCH)
466 .map(|d| d.as_secs()).unwrap_or(0);
467 let created_at = created_at_secs.min(now_secs);
469
470 let display_name = member_name.unwrap_or(member_npub);
471 let content = event_type.display_message(display_name);
472
473 let mut tags: Vec<Vec<String>> = vec![
474 vec!["d".to_string(), "system-event".to_string()],
475 vec!["event-type".to_string(), event_type.as_u8().to_string()],
476 vec!["member".to_string(), member_npub.to_string()],
477 ];
478 if let Some(by) = invited_by {
479 tags.push(vec!["invited-by".to_string(), by.to_string()]);
480 if let Some(l) = invited_label.filter(|l| !l.is_empty()) {
481 tags.push(vec!["invited-label".to_string(), l.to_string()]);
482 }
483 }
484 let tags_json = serde_json::to_string(&tags)
485 .map_err(|e| format!("Failed to serialize tags: {}", e))?;
486
487 let conn = super::get_write_connection_guard_static()?;
488 let rows = conn.execute(
489 r#"INSERT OR IGNORE INTO events (
490 id, kind, chat_id, user_id, content, tags, reference_id,
491 created_at, received_at, mine, pending, failed, wrapper_event_id, npub
492 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14)"#,
493 rusqlite::params![
494 event_id,
495 event_kind::APPLICATION_SPECIFIC as i32,
496 chat_id, None::<i64>, content, tags_json, None::<String>,
497 created_at as i64, now_secs as i64,
498 0, 0, 0, None::<String>, member_npub,
499 ],
500 ).map_err(|e| format!("Failed to save system event: {}", e))?;
501
502 Ok(rows > 0)
503}
504
505pub async fn save_edit_event(
507 edit_id: &str,
508 message_id: &str,
509 new_content: &str,
510 emoji_tags: &[crate::types::EmojiTag],
511 chat_id: i64,
512 user_id: Option<i64>,
513 npub: &str,
514) -> Result<(), String> {
515 let now = std::time::SystemTime::now()
516 .duration_since(std::time::UNIX_EPOCH).unwrap();
517
518 let mut tags = vec![
521 vec!["e".to_string(), message_id.to_string(), "".to_string(), "edit".to_string()],
522 ];
523 for et in emoji_tags {
524 tags.push(vec!["emoji".to_string(), et.shortcode.clone(), et.url.clone()]);
525 }
526
527 let event = StoredEvent {
528 id: edit_id.to_string(),
529 kind: event_kind::MESSAGE_EDIT,
530 chat_id,
531 user_id,
532 content: new_content.to_string(),
533 tags,
534 reference_id: Some(message_id.to_string()),
535 created_at: now.as_secs(),
536 received_at: now.as_millis() as u64,
537 mine: true,
538 pending: false,
539 failed: false,
540 wrapper_event_id: None,
541 npub: Some(npub.to_string()),
542 preview_metadata: None,
543 };
544
545 save_event(&event).await
546}
547
548pub async fn delete_event(event_id: &str) -> Result<(), String> {
550 let conn = super::get_write_connection_guard_static()?;
551 if let Ok(Some((chat_row, at))) = conn.query_row(
555 "SELECT chat_id, created_at FROM events WHERE id = ?1",
556 rusqlite::params![event_id],
557 |r| Ok((r.get::<_, i64>(0)?, r.get::<_, i64>(1)?)),
558 ).optional() {
559 conn.execute(
560 "UPDATE chats SET last_read = COALESCE(( \
561 SELECT id FROM events WHERE chat_id = ?1 AND id != ?2 AND created_at <= ?3 \
562 ORDER BY created_at DESC, id DESC LIMIT 1), '') \
563 WHERE id = ?1 AND last_read = ?2",
564 rusqlite::params![chat_row, event_id, at],
565 ).map_err(|e| format!("read-marker retreat: {e}"))?;
566 }
567 conn.execute(
568 "DELETE FROM events WHERE id = ?1",
569 rusqlite::params![event_id],
570 ).map_err(|e| format!("Failed to delete event: {}", e))?;
571 Ok(())
572}
573
574pub fn event_author(event_id: &str) -> Result<Option<String>, String> {
577 let conn = match super::get_db_connection_guard_static() {
578 Ok(c) => c,
579 Err(_) => return Ok(None),
580 };
581 conn.query_row(
582 "SELECT npub FROM events WHERE id = ?1",
583 rusqlite::params![event_id],
584 |row| row.get::<_, Option<String>>(0),
585 )
586 .optional()
587 .map(|o| o.flatten())
588 .map_err(|e| format!("Failed to read event author: {}", e))
589}
590
591pub fn event_delete_context(event_id: &str) -> Result<Option<(String, bool, Option<String>)>, String> {
595 let conn = match super::get_db_connection_guard_static() {
596 Ok(c) => c,
597 Err(_) => return Ok(None),
598 };
599 conn.query_row(
600 "SELECT c.chat_identifier, e.mine, e.npub \
601 FROM events e JOIN chats c ON c.id = e.chat_id \
602 WHERE e.id = ?1",
603 rusqlite::params![event_id],
604 |row| {
605 Ok((
606 row.get::<_, String>(0)?,
607 row.get::<_, i32>(1)? != 0,
608 row.get::<_, Option<String>>(2)?,
609 ))
610 },
611 )
612 .optional()
613 .map_err(|e| format!("Failed to read event delete context: {}", e))
614}
615
616pub fn message_exists_in_db(message_id: &str) -> Result<bool, String> {
618 let conn = match super::get_db_connection_guard_static() {
619 Ok(c) => c,
620 Err(_) => return Ok(false),
621 };
622 conn.query_row(
623 "SELECT EXISTS(SELECT 1 FROM events WHERE id = ?1)",
624 rusqlite::params![message_id],
625 |row| row.get(0),
626 ).map_err(|e| format!("Failed to check event existence: {}", e))
627}
628
629pub fn wrapper_event_exists(wrapper_event_id: &str) -> Result<bool, String> {
631 let conn = match super::get_db_connection_guard_static() {
632 Ok(c) => c,
633 Err(_) => return Ok(false),
634 };
635 conn.query_row(
636 "SELECT EXISTS(SELECT 1 FROM events WHERE wrapper_event_id = ?1)",
637 rusqlite::params![wrapper_event_id],
638 |row| row.get(0),
639 ).map_err(|e| format!("Failed to check wrapper event existence: {}", e))
640}
641
642pub fn update_wrapper_event_id(event_id: &str, wrapper_event_id: &str) -> Result<bool, String> {
645 let conn = match super::get_write_connection_guard_static() {
646 Ok(c) => c,
647 Err(_) => return Ok(false),
648 };
649 let rows = conn.execute(
650 "UPDATE events SET wrapper_event_id = ?1 WHERE id = ?2 AND (wrapper_event_id IS NULL OR wrapper_event_id = '')",
651 rusqlite::params![wrapper_event_id, event_id],
652 ).map_err(|e| format!("Failed to update wrapper event ID: {}", e))?;
653 Ok(rows > 0)
654}
655
656pub fn get_chat_message_count(chat_id: i64) -> Result<usize, String> {
658 let conn = super::get_db_connection_guard_static()?;
659 let count: i64 = conn.query_row(
664 &format!(
665 "SELECT COUNT(*) FROM events WHERE chat_id = ?1 AND kind IN ({}, {}, {})",
666 event_kind::CHAT_MESSAGE, event_kind::PRIVATE_DIRECT_MESSAGE, event_kind::FILE_ATTACHMENT
667 ),
668 rusqlite::params![chat_id],
669 |row| row.get(0),
670 ).map_err(|e| format!("Failed to count messages: {}", e))?;
671 Ok(count as usize)
672}
673
674pub fn get_pivx_payments_for_chat(conversation_id: &str) -> Result<Vec<StoredEvent>, String> {
676 let conn = super::get_db_connection_guard_static()?;
677 let chat_id: i64 = conn.query_row(
678 "SELECT id FROM chats WHERE chat_identifier = ?1",
679 rusqlite::params![conversation_id], |row| row.get(0)
680 ).map_err(|_| "Chat not found")?;
681
682 let mut stmt = conn.prepare(
683 "SELECT id, kind, chat_id, user_id, content, tags, reference_id, \
684 created_at, received_at, mine, pending, failed, wrapper_event_id, npub \
685 FROM events WHERE chat_id = ?1 AND kind = ?2 ORDER BY created_at ASC, received_at ASC"
686 ).map_err(|e| format!("Failed to prepare: {}", e))?;
687
688 let rows = stmt.query_map(
689 rusqlite::params![chat_id, event_kind::APPLICATION_SPECIFIC as i32],
690 |row| {
691 let tags_json: String = row.get(5)?;
692 let tags: Vec<Vec<String>> = serde_json::from_str(&tags_json).unwrap_or_default();
693 Ok(StoredEvent {
694 id: row.get(0)?, kind: row.get::<_, i32>(1)? as u16,
695 chat_id: row.get(2)?, user_id: row.get(3)?, content: row.get(4)?,
696 tags, reference_id: row.get(6)?,
697 created_at: row.get::<_, i64>(7)? as u64, received_at: row.get::<_, i64>(8)? as u64,
698 mine: row.get::<_, i32>(9)? != 0, pending: row.get::<_, i32>(10)? != 0,
699 failed: row.get::<_, i32>(11)? != 0, wrapper_event_id: row.get(12)?,
700 npub: row.get(13)?, preview_metadata: None,
701 })
702 }
703 ).map_err(|e| format!("Failed to query: {}", e))?;
704
705 let mut payments = Vec::new();
706 for row in rows {
707 let event = row.map_err(|e| format!("Failed to read event: {}", e))?;
708 if event.tags.iter().any(|t| t.len() >= 2 && t[0] == "d" && t[1] == "pivx-payment") {
709 payments.push(event);
710 }
711 }
712 Ok(payments)
713}
714
715pub fn get_system_events_for_chat(conversation_id: &str) -> Result<Vec<StoredEvent>, String> {
717 let conn = super::get_db_connection_guard_static()?;
718 let chat_id: i64 = conn.query_row(
719 "SELECT id FROM chats WHERE chat_identifier = ?1",
720 rusqlite::params![conversation_id], |row| row.get(0)
721 ).map_err(|_| "Chat not found")?;
722
723 let mut stmt = conn.prepare(
724 "SELECT id, kind, chat_id, user_id, content, tags, reference_id, \
725 created_at, received_at, mine, pending, failed, wrapper_event_id, npub \
726 FROM events WHERE chat_id = ?1 AND kind = ?2 ORDER BY created_at ASC, received_at ASC"
727 ).map_err(|e| format!("Failed to prepare: {}", e))?;
728
729 let rows = stmt.query_map(
730 rusqlite::params![chat_id, event_kind::APPLICATION_SPECIFIC as i32],
731 |row| {
732 let tags_json: String = row.get(5)?;
733 let tags: Vec<Vec<String>> = serde_json::from_str(&tags_json).unwrap_or_default();
734 Ok(StoredEvent {
735 id: row.get(0)?, kind: row.get::<_, i32>(1)? as u16,
736 chat_id: row.get(2)?, user_id: row.get(3)?, content: row.get(4)?,
737 tags, reference_id: row.get(6)?,
738 created_at: row.get::<_, i64>(7)? as u64, received_at: row.get::<_, i64>(8)? as u64,
739 mine: row.get::<_, i32>(9)? != 0, pending: row.get::<_, i32>(10)? != 0,
740 failed: row.get::<_, i32>(11)? != 0, wrapper_event_id: row.get(12)?,
741 npub: row.get(13)?, preview_metadata: None,
742 })
743 }
744 ).map_err(|e| format!("Failed to query: {}", e))?;
745
746 let mut events = Vec::new();
747 for row in rows {
748 let event = row.map_err(|e| format!("Failed to read event: {}", e))?;
749 if event.tags.iter().any(|t| t.len() >= 2 && t[0] == "d" && t[1] == "system-event") {
750 events.push(event);
751 }
752 }
753 Ok(events)
754}
755
756fn parse_event_row(row: &rusqlite::Row) -> rusqlite::Result<StoredEvent> {
762 let tags_json: String = row.get(5)?;
763 let tags: Vec<Vec<String>> = serde_json::from_str(&tags_json).unwrap_or_default();
764
765 Ok(StoredEvent {
766 id: row.get(0)?,
767 kind: row.get::<_, i32>(1)? as u16,
768 chat_id: row.get(2)?,
769 user_id: row.get(3)?,
770 content: row.get(4)?,
771 tags,
772 reference_id: row.get(6)?,
773 created_at: row.get::<_, i64>(7)? as u64,
774 received_at: row.get::<_, i64>(8)? as u64,
775 mine: row.get::<_, i32>(9)? != 0,
776 pending: row.get::<_, i32>(10)? != 0,
777 failed: row.get::<_, i32>(11)? != 0,
778 wrapper_event_id: row.get(12)?,
779 npub: row.get(13)?,
780 preview_metadata: row.get(14)?,
781 })
782}
783
784pub async fn get_events(
787 chat_id: i64,
788 kinds: Option<&[u16]>,
789 limit: usize,
790 offset: usize,
791) -> Result<Vec<StoredEvent>, String> {
792 let events: Vec<StoredEvent> = {
793 let conn = super::get_db_connection_guard_static()?;
794
795 if let Some(k) = kinds {
796 let kind_placeholders: String = (0..k.len())
797 .map(|i| format!("?{}", i + 2))
798 .collect::<Vec<_>>()
799 .join(",");
800 let limit_param = k.len() + 2;
801 let offset_param = k.len() + 3;
802
803 let sql = format!(
804 "SELECT id, kind, chat_id, user_id, content, tags, reference_id, \
805 created_at, received_at, mine, pending, failed, wrapper_event_id, npub, preview_metadata \
806 FROM events WHERE chat_id = ?1 AND kind IN ({}) \
807 ORDER BY created_at DESC, received_at DESC \
808 LIMIT ?{} OFFSET ?{}",
809 kind_placeholders, limit_param, offset_param
810 );
811
812 let mut stmt = conn.prepare(&sql)
813 .map_err(|e| format!("Failed to prepare events query: {}", e))?;
814
815 match k.len() {
816 1 => {
817 let rows = stmt.query_map(
818 rusqlite::params![chat_id, k[0] as i32, limit as i64, offset as i64],
819 parse_event_row
820 ).map_err(|e| format!("Failed to query events: {}", e))?;
821 rows.filter_map(|r| r.ok()).collect()
822 },
823 2 => {
824 let rows = stmt.query_map(
825 rusqlite::params![chat_id, k[0] as i32, k[1] as i32, limit as i64, offset as i64],
826 parse_event_row
827 ).map_err(|e| format!("Failed to query events: {}", e))?;
828 rows.filter_map(|r| r.ok()).collect()
829 },
830 3 => {
831 let rows = stmt.query_map(
832 rusqlite::params![chat_id, k[0] as i32, k[1] as i32, k[2] as i32, limit as i64, offset as i64],
833 parse_event_row
834 ).map_err(|e| format!("Failed to query events: {}", e))?;
835 rows.filter_map(|r| r.ok()).collect()
836 },
837 _ => return Err("Unsupported number of kinds".to_string()),
838 }
839 } else {
840 let mut stmt = conn.prepare(
841 "SELECT id, kind, chat_id, user_id, content, tags, reference_id, \
842 created_at, received_at, mine, pending, failed, wrapper_event_id, npub, preview_metadata \
843 FROM events WHERE chat_id = ?1 \
844 ORDER BY created_at DESC, received_at DESC \
845 LIMIT ?2 OFFSET ?3"
846 ).map_err(|e| format!("Failed to prepare events query: {}", e))?;
847
848 let rows = stmt.query_map(
849 rusqlite::params![chat_id, limit as i64, offset as i64],
850 parse_event_row
851 ).map_err(|e| format!("Failed to query events: {}", e))?;
852 rows.filter_map(|r| r.ok()).collect()
853 }
854 };
855
856 let mut decrypted = Vec::with_capacity(events.len());
858 for mut event in events {
859 if event.kind == event_kind::CHAT_MESSAGE || event.kind == event_kind::PRIVATE_DIRECT_MESSAGE {
860 event.content = crate::crypto::maybe_decrypt(event.content).await
861 .unwrap_or_else(|_| "[Decryption failed]".to_string());
862 }
863 decrypted.push(event);
864 }
865
866 Ok(decrypted)
867}
868
869pub async fn get_related_events(
871 reference_ids: &[String],
872) -> Result<Vec<StoredEvent>, String> {
873 if reference_ids.is_empty() {
874 return Ok(Vec::new());
875 }
876
877 let conn = super::get_db_connection_guard_static()?;
878
879 let placeholders: String = reference_ids.iter().map(|_| "?").collect::<Vec<_>>().join(",");
880 let sql = format!(
881 "SELECT id, kind, chat_id, user_id, content, tags, reference_id, \
882 created_at, received_at, mine, pending, failed, wrapper_event_id, npub, preview_metadata \
883 FROM events WHERE reference_id IN ({}) \
884 ORDER BY created_at ASC, received_at ASC",
885 placeholders
886 );
887
888 let mut stmt = conn.prepare(&sql)
889 .map_err(|e| format!("Failed to prepare related events query: {}", e))?;
890
891 let params: Vec<&dyn rusqlite::ToSql> = reference_ids.iter()
892 .map(|s| s as &dyn rusqlite::ToSql)
893 .collect();
894
895 let events: Vec<StoredEvent> = stmt.query_map(params.as_slice(), parse_event_row)
896 .map_err(|e| format!("Failed to query related events: {}", e))?
897 .filter_map(|r| r.ok())
898 .collect();
899
900 Ok(events)
901}
902
903pub struct ReplyContext {
905 pub content: String,
906 pub npub: Option<String>,
907 pub has_attachment: bool,
908 pub extension: Option<String>,
911}
912
913pub async fn get_reply_contexts(
915 message_ids: &[String],
916) -> Result<std::collections::HashMap<String, ReplyContext>, String> {
917 use std::collections::HashMap;
918
919 if message_ids.is_empty() {
920 return Ok(HashMap::new());
921 }
922
923 let (events, edits): (Vec<(String, i32, String, Option<String>, Option<String>)>, Vec<(String, String)>) = {
924 let conn = super::get_db_connection_guard_static()?;
925
926 let placeholders: String = (0..message_ids.len())
927 .map(|i| format!("?{}", i + 1))
928 .collect::<Vec<_>>()
929 .join(",");
930
931 let sql = format!(
933 "SELECT id, kind, content, npub, tags FROM events WHERE id IN ({})",
934 placeholders
935 );
936 let mut stmt = conn.prepare(&sql)
937 .map_err(|e| format!("Failed to prepare reply context query: {}", e))?;
938
939 let params: Vec<&str> = message_ids.iter().map(|s| s.as_str()).collect();
940 let params_dyn: Vec<&dyn rusqlite::ToSql> = params.iter().map(|s| s as &dyn rusqlite::ToSql).collect();
941
942 let rows = stmt.query_map(params_dyn.as_slice(), |row| {
943 Ok((row.get::<_, String>(0)?, row.get::<_, i32>(1)?,
944 row.get::<_, String>(2)?, row.get::<_, Option<String>>(3)?,
945 row.get::<_, Option<String>>(4)?))
946 }).map_err(|e| format!("Failed to query reply contexts: {}", e))?;
947 let events_result: Vec<_> = rows.filter_map(|r| r.ok()).collect();
948 drop(stmt);
949
950 let edit_sql = format!(
952 "SELECT reference_id, content FROM events \
953 WHERE kind = {} AND reference_id IN ({}) \
954 ORDER BY created_at DESC, received_at DESC",
955 event_kind::MESSAGE_EDIT, placeholders
956 );
957 let mut edit_stmt = conn.prepare(&edit_sql)
958 .map_err(|e| format!("Failed to prepare edit query: {}", e))?;
959 let edit_rows = edit_stmt.query_map(params_dyn.as_slice(), |row| {
960 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
961 }).map_err(|e| format!("Failed to query edits: {}", e))?;
962 let edits_result: Vec<_> = edit_rows.filter_map(|r| r.ok()).collect();
963
964 (events_result, edits_result)
965 };
966
967 let mut latest_edits: HashMap<String, String> = HashMap::new();
969 for (ref_id, content) in edits {
970 latest_edits.entry(ref_id).or_insert(content);
971 }
972
973 let file_ids: Vec<String> = events.iter()
975 .filter(|(_, kind, _, _, _)| *kind == event_kind::FILE_ATTACHMENT as i32)
976 .map(|(id, _, _, _, _)| id.clone())
977 .collect();
978 let atts_by_event = super::attachments::get_attachments_for_events(&file_ids).unwrap_or_default();
979
980 let mut contexts = HashMap::new();
982 for (id, kind, original_content, npub, tags) in events {
983 let has_attachment = kind == event_kind::FILE_ATTACHMENT as i32;
984 let content_to_decrypt = latest_edits.get(&id).cloned().unwrap_or(original_content);
985
986 let decrypted_content = if kind == event_kind::CHAT_MESSAGE as i32
987 || kind == event_kind::PRIVATE_DIRECT_MESSAGE as i32
988 {
989 crate::crypto::maybe_decrypt(content_to_decrypt).await
990 .unwrap_or_else(|_| "[Decryption failed]".to_string())
991 } else {
992 String::new()
993 };
994
995 let extension = if has_attachment {
998 atts_by_event.get(&id)
999 .and_then(|atts| atts.first())
1000 .map(|a| a.extension.to_lowercase())
1001 .filter(|e| !e.is_empty())
1002 .or_else(|| tags.as_deref()
1003 .and_then(|t| serde_json::from_str::<Vec<Vec<String>>>(t).ok())
1004 .and_then(|parsed| parsed.into_iter()
1005 .find(|t| t.first().map(|k| k == "attachments").unwrap_or(false))
1006 .and_then(|t| t.into_iter().nth(1)))
1007 .and_then(|json| serde_json::from_str::<Vec<serde_json::Value>>(&json).ok())
1008 .and_then(|atts| atts.into_iter().next())
1009 .and_then(|a| a.get("extension").and_then(|e| e.as_str()).map(str::to_lowercase))
1010 .filter(|e| !e.is_empty()))
1011 } else {
1012 None
1013 };
1014
1015 contexts.insert(id, ReplyContext { content: decrypted_content, npub, has_attachment, extension });
1016 }
1017
1018 Ok(contexts)
1019}
1020
1021pub async fn populate_reply_contexts(messages: Vec<&mut Message>) -> Result<(), String> {
1028 let ids: Vec<String> = messages
1029 .iter()
1030 .filter(|m| !m.replied_to.is_empty())
1031 .map(|m| m.replied_to.clone())
1032 .collect();
1033 if ids.is_empty() {
1034 return Ok(());
1035 }
1036 let contexts = get_reply_contexts(&ids).await?;
1037 for message in messages {
1038 if let Some(ctx) = contexts.get(&message.replied_to) {
1039 message.replied_to_content = Some(ctx.content.clone());
1040 message.replied_to_npub = ctx.npub.clone();
1041 message.replied_to_has_attachment = Some(ctx.has_attachment);
1042 message.replied_to_attachment_extension = ctx.extension.clone();
1043 }
1044 }
1045 Ok(())
1046}
1047
1048pub async fn populate_reply_context(message: &mut Message) -> Result<(), String> {
1051 if message.replied_to.is_empty() {
1052 return Ok(());
1053 }
1054
1055 let contexts = get_reply_contexts(&[message.replied_to.clone()]).await?;
1056
1057 if let Some(ctx) = contexts.get(&message.replied_to) {
1058 message.replied_to_content = Some(ctx.content.clone());
1059 message.replied_to_npub = ctx.npub.clone();
1060 message.replied_to_has_attachment = Some(ctx.has_attachment);
1061 message.replied_to_attachment_extension = ctx.extension.clone();
1062 }
1063
1064 Ok(())
1065}
1066
1067pub fn is_own_event(event_id: &str) -> bool {
1071 let Ok(conn) = super::get_db_connection_guard_static() else {
1072 return false;
1073 };
1074 conn.query_row(
1075 "SELECT mine FROM events WHERE id = ?1",
1076 [event_id],
1077 |row| row.get::<_, i64>(0),
1078 )
1079 .map(|mine| mine == 1)
1080 .unwrap_or(false)
1081}
1082
1083fn extract_tag_from_json(tags_json: &str, key: &str) -> Option<String> {
1089 if tags_json.len() <= 2 { return None; }
1090 let pattern = format!("[\"{}\"", key);
1091 if !tags_json.contains(&pattern) { return None; }
1092 let tags: Vec<Vec<String>> = serde_json::from_str(tags_json).ok()?;
1093 tags.into_iter()
1094 .find(|tag| tag.first().map(|s| s.as_str()) == Some(key))
1095 .and_then(|tag| tag.into_iter().nth(1))
1096}
1097
1098
1099fn normalize_reaction_author(author: String) -> String {
1103 if author.len() == 64 && author.bytes().all(|b| b.is_ascii_hexdigit()) {
1104 if let Ok(pk) = nostr_sdk::prelude::PublicKey::from_hex(&author) {
1105 use nostr_sdk::prelude::ToBech32;
1106 let Ok(npub) = pk.to_bech32();
1107 return npub;
1108 }
1109 }
1110 author
1111}
1112fn extract_reaction_emoji_url(tags: &[Vec<String>], content: &str) -> Option<String> {
1118 if !content.starts_with(':') || !content.ends_with(':') || content.len() < 3 {
1119 return None;
1120 }
1121 let sc = &content[1..content.len() - 1];
1122 tags.iter().find_map(|t| {
1123 if t.len() >= 3 && t[0] == "emoji" && t[1] == sc {
1124 Some(t[2].clone())
1125 } else {
1126 None
1127 }
1128 })
1129}
1130
1131fn extract_reply_tag_from_json(tags_json: &str) -> Option<String> {
1133 if tags_json.len() <= 2 { return None; }
1134 if !tags_json.contains("[\"e\"") { return None; }
1135 let tags: Vec<Vec<String>> = serde_json::from_str(tags_json).ok()?;
1136 tags.into_iter()
1137 .find(|tag| {
1138 tag.first().map(|s| s.as_str()) == Some("e")
1139 && tag.get(3).map(|s| s.as_str()) == Some("reply")
1140 })
1141 .and_then(|tag| tag.into_iter().nth(1))
1142}
1143
1144pub async fn get_message_views(
1149 chat_id: i64,
1150 limit: usize,
1151 offset: usize,
1152) -> Result<Vec<Message>, String> {
1153 let message_kinds = [event_kind::CHAT_MESSAGE, event_kind::PRIVATE_DIRECT_MESSAGE, event_kind::FILE_ATTACHMENT];
1155 let message_events = get_events(chat_id, Some(&message_kinds), limit, offset).await?;
1156
1157 compose_message_views(message_events).await
1158}
1159
1160async fn compose_message_views(message_events: Vec<StoredEvent>) -> Result<Vec<Message>, String> {
1165 use std::collections::HashMap;
1166
1167 if message_events.is_empty() {
1168 return Ok(Vec::new());
1169 }
1170
1171 let message_ids: Vec<String> = message_events.iter().map(|e| e.id.clone()).collect();
1173 let related_events = get_related_events(&message_ids).await?;
1174
1175 let mut reactions_by_msg: HashMap<String, Vec<Reaction>> = HashMap::new();
1176 let mut edits_by_msg: HashMap<String, Vec<(u64, String, Vec<crate::types::EmojiTag>)>> = HashMap::new();
1177
1178 for event in related_events {
1179 if let Some(ref_id) = &event.reference_id {
1180 match event.kind {
1181 k if k == event_kind::REACTION => {
1182 let emoji_url = extract_reaction_emoji_url(&event.tags, &event.content);
1183 reactions_by_msg.entry(ref_id.clone()).or_default().push(Reaction {
1184 id: event.id.clone(),
1185 reference_id: ref_id.clone(),
1186 author_id: normalize_reaction_author(event.npub.clone().unwrap_or_default()),
1187 emoji: event.content.clone(),
1188 emoji_url,
1189 });
1190 }
1191 k if k == event_kind::MESSAGE_EDIT => {
1192 let decrypted = crate::crypto::maybe_decrypt(event.content.clone()).await
1193 .unwrap_or_else(|_| event.content.clone());
1194 let edit_emoji = crate::types::EmojiTag::extract_from_stored(&event.tags);
1195 edits_by_msg.entry(ref_id.clone()).or_default().push((event.created_at * 1000, decrypted, edit_emoji));
1196 }
1197 _ => {}
1198 }
1199 }
1200 }
1201
1202 for edits in edits_by_msg.values_mut() {
1203 edits.sort_by_key(|(ts, _, _)| *ts);
1204 }
1205
1206 let attach_ids: Vec<String> = message_events.iter()
1209 .filter(|e| e.kind == event_kind::FILE_ATTACHMENT || e.kind == event_kind::CHAT_MESSAGE)
1210 .map(|e| e.id.clone())
1211 .collect();
1212 let mut attachments_by_msg = super::attachments::get_attachments_for_events(&attach_ids)
1213 .unwrap_or_default();
1214 for event in &message_events {
1215 if event.kind != event_kind::FILE_ATTACHMENT && event.kind != event_kind::CHAT_MESSAGE {
1216 continue;
1217 }
1218 if attachments_by_msg.contains_key(&event.id) {
1219 continue;
1220 }
1221 if let Some(json) = event.get_tag("attachments") {
1222 if let Ok(atts) = serde_json::from_str::<Vec<Attachment>>(json) {
1223 if !atts.is_empty() {
1224 attachments_by_msg.insert(event.id.clone(), atts);
1225 }
1226 }
1227 }
1228 }
1229
1230 let mut messages = Vec::with_capacity(message_events.len());
1232 for event in message_events {
1233 let replied_to = event.get_reply_reference().unwrap_or("").to_string();
1234 let at = event.timestamp_ms();
1235 let reactions = reactions_by_msg.remove(&event.id).unwrap_or_default();
1236 let attachments = attachments_by_msg.remove(&event.id).unwrap_or_default();
1237
1238 let original_content = if event.kind == event_kind::FILE_ATTACHMENT {
1239 String::new()
1240 } else {
1241 event.content.clone()
1242 };
1243
1244 let original_emoji = crate::types::EmojiTag::extract_from_stored(&event.tags);
1247 let (content, edited, edit_history, emoji_tags) = if let Some(edits) = edits_by_msg.remove(&event.id) {
1248 let mut history = Vec::with_capacity(edits.len() + 1);
1249 history.push(crate::types::EditEntry { content: original_content.clone(), edited_at: at });
1250 for (ts, c, _) in &edits {
1251 history.push(crate::types::EditEntry { content: c.clone(), edited_at: *ts });
1252 }
1253 let (latest, latest_emoji) = edits.last()
1254 .map(|(_, c, e)| (c.clone(), e.clone()))
1255 .unwrap_or_else(|| (original_content.clone(), original_emoji.clone()));
1256 (latest, true, Some(history), latest_emoji)
1257 } else {
1258 (original_content, false, None, original_emoji)
1259 };
1260
1261 let preview_metadata = event.preview_metadata
1262 .and_then(|json| serde_json::from_str(&json).ok());
1263
1264 let addressed_bots = extract_bot_tags(&event.tags);
1265 let expiration = extract_expiration_tag(&event.tags);
1266 messages.push(Message {
1267 expiration,
1268 id: event.id, content, replied_to,
1269 replied_to_content: None, replied_to_npub: None, replied_to_has_attachment: None,
1270 replied_to_attachment_extension: None,
1271 preview_metadata, attachments, reactions, at,
1272 pending: event.pending, failed: event.failed, mine: event.mine,
1273 npub: event.npub, wrapper_event_id: event.wrapper_event_id,
1274 edited, edit_history,
1275 emoji_tags,
1276 addressed_bots,
1277 });
1278 }
1279
1280 crate::self_destruct::strip_expired(&mut messages);
1284
1285 let reply_ids: Vec<String> = messages.iter()
1287 .filter(|m| !m.replied_to.is_empty())
1288 .map(|m| m.replied_to.clone())
1289 .collect();
1290
1291 if !reply_ids.is_empty() {
1292 let contexts = get_reply_contexts(&reply_ids).await?;
1293 for msg in &mut messages {
1294 if let Some(ctx) = contexts.get(&msg.replied_to) {
1295 msg.replied_to_content = Some(ctx.content.clone());
1296 msg.replied_to_npub = ctx.npub.clone();
1297 msg.replied_to_has_attachment = Some(ctx.has_attachment);
1298 msg.replied_to_attachment_extension = ctx.extension.clone();
1299 }
1300 }
1301 }
1302
1303 Ok(messages)
1304}
1305
1306pub async fn get_messages_around(
1314 chat_id: i64,
1315 anchor_id: &str,
1316 before: usize,
1317 after: usize,
1318) -> Result<Vec<Message>, String> {
1319 let message_kinds = [event_kind::CHAT_MESSAGE, event_kind::PRIVATE_DIRECT_MESSAGE, event_kind::FILE_ATTACHMENT];
1320
1321 let message_events: Vec<StoredEvent> = {
1322 let conn = super::get_db_connection_guard_static()?;
1323
1324 let (anchor_at, anchor_rt, anchor_rowid): (i64, i64, i64) = conn.query_row(
1330 "SELECT created_at, received_at, rowid FROM events WHERE id = ?1",
1331 rusqlite::params![anchor_id],
1332 |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
1333 ).map_err(|e| format!("Anchor message not found: {}", e))?;
1334
1335 let kind_placeholders: String = (0..message_kinds.len())
1337 .map(|i| format!("?{}", i + 2))
1338 .collect::<Vec<_>>()
1339 .join(",");
1340 let cols = "id, kind, chat_id, user_id, content, tags, reference_id, \
1341 created_at, received_at, mine, pending, failed, wrapper_event_id, npub, preview_metadata";
1342
1343 let older_sql = format!(
1345 "SELECT {} FROM events WHERE chat_id = ?1 AND kind IN ({}) \
1346 AND (created_at < ?5 OR (created_at = ?5 AND (received_at < ?6 \
1347 OR (received_at = ?6 AND rowid <= ?7)))) \
1348 ORDER BY created_at DESC, received_at DESC, rowid DESC LIMIT ?8",
1349 cols, kind_placeholders
1350 );
1351 let mut older_stmt = conn.prepare(&older_sql)
1352 .map_err(|e| format!("Failed to prepare older window query: {}", e))?;
1353 let older_rows = older_stmt.query_map(
1354 rusqlite::params![
1355 chat_id,
1356 message_kinds[0] as i32, message_kinds[1] as i32, message_kinds[2] as i32,
1357 anchor_at, anchor_rt, anchor_rowid, before as i64
1358 ],
1359 parse_event_row,
1360 ).map_err(|e| format!("Failed to query older window: {}", e))?;
1361 let mut older: Vec<StoredEvent> = older_rows.filter_map(|r| r.ok()).collect();
1362 older.reverse(); let newer_sql = format!(
1366 "SELECT {} FROM events WHERE chat_id = ?1 AND kind IN ({}) \
1367 AND (created_at > ?5 OR (created_at = ?5 AND (received_at > ?6 \
1368 OR (received_at = ?6 AND rowid > ?7)))) \
1369 ORDER BY created_at ASC, received_at ASC, rowid ASC LIMIT ?8",
1370 cols, kind_placeholders
1371 );
1372 let mut newer_stmt = conn.prepare(&newer_sql)
1373 .map_err(|e| format!("Failed to prepare newer window query: {}", e))?;
1374 let newer_rows = newer_stmt.query_map(
1375 rusqlite::params![
1376 chat_id,
1377 message_kinds[0] as i32, message_kinds[1] as i32, message_kinds[2] as i32,
1378 anchor_at, anchor_rt, anchor_rowid, after as i64
1379 ],
1380 parse_event_row,
1381 ).map_err(|e| format!("Failed to query newer window: {}", e))?;
1382 let newer: Vec<StoredEvent> = newer_rows.filter_map(|r| r.ok()).collect();
1383
1384 older.into_iter().chain(newer).collect()
1385 };
1386
1387 let mut decrypted = Vec::with_capacity(message_events.len());
1389 for mut event in message_events {
1390 if event.kind == event_kind::CHAT_MESSAGE || event.kind == event_kind::PRIVATE_DIRECT_MESSAGE {
1391 event.content = crate::crypto::maybe_decrypt(event.content).await
1392 .unwrap_or_else(|_| "[Decryption failed]".to_string());
1393 }
1394 decrypted.push(event);
1395 }
1396
1397 compose_message_views(decrypted).await
1398}
1399
1400pub async fn get_all_chats_last_messages() -> Result<std::collections::HashMap<String, Vec<Message>>, String> {
1403 use std::collections::HashMap;
1404
1405 let message_events: Vec<(String, StoredEvent, String)> = {
1407 let conn = super::get_db_connection_guard_static()?;
1408 let mut stmt = conn.prepare(
1409 "SELECT c.chat_identifier, \
1410 e.id, e.kind, e.chat_id, e.user_id, e.content, e.tags, e.reference_id, \
1411 e.created_at, e.received_at, e.mine, e.pending, e.failed, e.wrapper_event_id, e.npub, e.preview_metadata \
1412 FROM chats c JOIN events e ON e.rowid = ( \
1413 SELECT e2.rowid FROM events e2 WHERE e2.chat_id = c.id \
1414 AND e2.kind IN (?1, ?2, ?3) \
1415 ORDER BY e2.created_at DESC, e2.received_at DESC LIMIT 1) \
1416 WHERE c.chat_type != 1"
1417 ).map_err(|e| format!("Failed to prepare: {}", e))?;
1418
1419 let rows = stmt.query_map(
1420 rusqlite::params![
1421 event_kind::CHAT_MESSAGE as i32,
1422 event_kind::PRIVATE_DIRECT_MESSAGE as i32,
1423 event_kind::FILE_ATTACHMENT as i32
1424 ],
1425 |row| {
1426 let chat_id: String = row.get(0)?;
1427 let tags_json: String = row.get(6)?;
1428 let event = StoredEvent {
1429 id: row.get(1)?, kind: row.get::<_, i32>(2)? as u16,
1430 chat_id: row.get(3)?, user_id: row.get(4)?, content: row.get(5)?,
1431 tags: Vec::new(), reference_id: row.get(7)?,
1433 created_at: row.get::<_, i64>(8)? as u64, received_at: row.get::<_, i64>(9)? as u64,
1434 mine: row.get::<_, i32>(10)? != 0, pending: row.get::<_, i32>(11)? != 0,
1435 failed: row.get::<_, i32>(12)? != 0, wrapper_event_id: row.get(13)?,
1436 npub: row.get(14)?, preview_metadata: row.get(15)?,
1437 };
1438 Ok((chat_id, event, tags_json))
1439 }
1440 ).map_err(|e| format!("Failed to query: {}", e))?;
1441 rows.filter_map(|r| r.ok()).collect()
1442 };
1443
1444 if message_events.is_empty() {
1445 return Ok(HashMap::new());
1446 }
1447
1448 let message_ids: Vec<String> = message_events.iter().map(|(_, e, _)| e.id.clone()).collect();
1450 let related_events = get_related_events(&message_ids).await?;
1451
1452 let mut reactions_by_msg: HashMap<String, Vec<Reaction>> = HashMap::new();
1453 let mut edits_by_msg: HashMap<String, Vec<(u64, String, Vec<crate::types::EmojiTag>)>> = HashMap::new();
1454
1455 for event in related_events {
1456 if let Some(ref_id) = &event.reference_id {
1457 match event.kind {
1458 k if k == event_kind::REACTION => {
1459 let emoji_url = extract_reaction_emoji_url(&event.tags, &event.content);
1460 reactions_by_msg.entry(ref_id.clone()).or_default().push(Reaction {
1461 id: event.id.clone(), reference_id: ref_id.clone(),
1462 author_id: normalize_reaction_author(event.npub.clone().unwrap_or_default()),
1463 emoji: event.content.clone(),
1464 emoji_url,
1465 });
1466 }
1467 k if k == event_kind::MESSAGE_EDIT => {
1468 let decrypted = crate::crypto::maybe_decrypt(event.content.clone()).await
1469 .unwrap_or_else(|_| event.content.clone());
1470 let edit_emoji = crate::types::EmojiTag::extract_from_stored(&event.tags);
1471 edits_by_msg.entry(ref_id.clone()).or_default().push((event.created_at * 1000, decrypted, edit_emoji));
1472 }
1473 _ => {}
1474 }
1475 }
1476 }
1477 for edits in edits_by_msg.values_mut() {
1478 edits.sort_by_key(|(ts, _, _)| *ts);
1479 }
1480
1481 let attach_ids: Vec<String> = message_events.iter()
1483 .filter(|(_, e, _)| e.kind == event_kind::FILE_ATTACHMENT || e.kind == event_kind::CHAT_MESSAGE)
1484 .map(|(_, e, _)| e.id.clone())
1485 .collect();
1486 let mut attachments_by_msg = super::attachments::get_attachments_for_events(&attach_ids)
1487 .unwrap_or_default();
1488 for (_, event, tags_json) in &message_events {
1489 if event.kind != event_kind::FILE_ATTACHMENT && event.kind != event_kind::CHAT_MESSAGE {
1490 continue;
1491 }
1492 if attachments_by_msg.contains_key(&event.id) {
1493 continue;
1494 }
1495 if let Some(val) = extract_tag_from_json(tags_json, "attachments") {
1496 if let Ok(atts) = serde_json::from_str::<Vec<Attachment>>(&val) {
1497 if !atts.is_empty() {
1498 attachments_by_msg.insert(event.id.clone(), atts);
1499 }
1500 }
1501 }
1502 }
1503
1504 let mut result: HashMap<String, Vec<Message>> = HashMap::new();
1506
1507 for (chat_identifier, event, tags_json) in message_events {
1508 let reactions = reactions_by_msg.remove(&event.id).unwrap_or_default();
1509 let attachments = attachments_by_msg.remove(&event.id).unwrap_or_default();
1510 let replied_to = extract_reply_tag_from_json(&tags_json).unwrap_or_default();
1511
1512 let original_content = if event.kind == event_kind::CHAT_MESSAGE
1514 || event.kind == event_kind::PRIVATE_DIRECT_MESSAGE
1515 {
1516 crate::crypto::maybe_decrypt(event.content.clone()).await
1517 .unwrap_or_else(|_| "[Decryption failed]".to_string())
1518 } else {
1519 String::new()
1520 };
1521
1522 let stored_tags = serde_json::from_str::<Vec<Vec<String>>>(&tags_json).unwrap_or_default();
1523 let original_emoji = crate::types::EmojiTag::extract_from_stored(&stored_tags);
1524 let addressed_bots = extract_bot_tags(&stored_tags);
1525 let expiration = extract_expiration_tag(&stored_tags);
1526 let (content, edited, edit_history, emoji_tags) = if let Some(edits) = edits_by_msg.remove(&event.id) {
1528 let (latest, latest_emoji) = edits.last()
1529 .map(|(_, c, e)| (c.clone(), e.clone()))
1530 .unwrap_or_else(|| (original_content.clone(), original_emoji.clone()));
1531 let history: Vec<crate::types::EditEntry> = std::iter::once(crate::types::EditEntry {
1532 content: original_content, edited_at: event.created_at * 1000,
1533 }).chain(edits.into_iter().map(|(ts, c, _)| crate::types::EditEntry { content: c, edited_at: ts }))
1534 .collect();
1535 (latest, true, Some(history), latest_emoji)
1536 } else {
1537 (original_content, false, None, original_emoji)
1538 };
1539
1540 let preview_metadata = event.preview_metadata
1541 .and_then(|json| serde_json::from_str(&json).ok());
1542
1543 result.entry(chat_identifier).or_default().push(Message {
1544 expiration,
1545 id: event.id, content, replied_to,
1546 replied_to_content: None, replied_to_npub: None, replied_to_has_attachment: None,
1547 replied_to_attachment_extension: None,
1548 preview_metadata, attachments, reactions, at: event.created_at * 1000,
1549 pending: event.pending, failed: event.failed, mine: event.mine,
1550 npub: event.npub, wrapper_event_id: event.wrapper_event_id,
1551 edited, edit_history,
1552 emoji_tags,
1553 addressed_bots,
1554 });
1555 }
1556
1557 let reply_ids: Vec<String> = result.values()
1561 .flatten()
1562 .filter(|m| !m.replied_to.is_empty())
1563 .map(|m| m.replied_to.clone())
1564 .collect();
1565
1566 if !reply_ids.is_empty() {
1567 let contexts = get_reply_contexts(&reply_ids).await?;
1568 for msg in result.values_mut().flatten() {
1569 if let Some(ctx) = contexts.get(&msg.replied_to) {
1570 msg.replied_to_content = Some(ctx.content.clone());
1571 msg.replied_to_npub = ctx.npub.clone();
1572 msg.replied_to_has_attachment = Some(ctx.has_attachment);
1573 msg.replied_to_attachment_extension = ctx.extension.clone();
1574 }
1575 }
1576 }
1577
1578 for msgs in result.values_mut() {
1582 crate::self_destruct::strip_expired(msgs);
1583 }
1584
1585 Ok(result)
1586}
1587
1588pub async fn unread_counts() -> Result<std::collections::HashMap<String, u32>, String> {
1595 let conn = super::get_db_connection_guard_static()?;
1596 let mut stmt = conn
1602 .prepare(
1603 "WITH anchors AS ( \
1604 SELECT c.id AS chat_id, c.chat_identifier AS chat_identifier, \
1605 COALESCE(MAX(e.created_at), 0) AS anchor_ts \
1606 FROM chats c \
1607 LEFT JOIN events e ON e.chat_id = c.id \
1608 AND ((e.mine = 1 AND e.kind IN (?1, ?2, ?3)) OR e.id = c.last_read) \
1609 GROUP BY c.id \
1610 ) \
1611 SELECT a.chat_identifier, COUNT(*) AS unread \
1612 FROM events e JOIN anchors a ON a.chat_id = e.chat_id \
1613 WHERE e.kind IN (?1, ?2, ?3) AND e.mine = 0 AND e.created_at > a.anchor_ts \
1614 GROUP BY a.chat_identifier",
1615 )
1616 .map_err(|e| format!("prepare unread_counts: {e}"))?;
1617 let rows = stmt
1618 .query_map(
1619 rusqlite::params![
1620 event_kind::CHAT_MESSAGE as i32,
1621 event_kind::PRIVATE_DIRECT_MESSAGE as i32,
1622 event_kind::FILE_ATTACHMENT as i32
1623 ],
1624 |row| Ok((row.get::<_, String>(0)?, row.get::<_, i64>(1)? as u32)),
1625 )
1626 .map_err(|e| format!("query unread_counts: {e}"))?;
1627 let mut out = std::collections::HashMap::new();
1628 for r in rows.flatten() {
1629 out.insert(r.0, r.1);
1630 }
1631 Ok(out)
1632}
1633
1634pub async fn unread_count_for_chat(chat_identifier: &str) -> Result<u32, String> {
1637 let conn = super::get_db_connection_guard_static()?;
1638 let count: i64 = conn
1639 .query_row(
1640 "SELECT COUNT(*) FROM events e JOIN chats c ON e.chat_id = c.id \
1641 WHERE c.chat_identifier = ?4 AND e.kind IN (?1, ?2, ?3) AND e.mine = 0 \
1642 AND e.created_at > COALESCE(( \
1643 SELECT MAX(e2.created_at) FROM events e2 \
1644 WHERE e2.chat_id = c.id \
1645 AND ((e2.mine = 1 AND e2.kind IN (?1, ?2, ?3)) OR e2.id = c.last_read)), 0)",
1646 rusqlite::params![
1647 event_kind::CHAT_MESSAGE as i32,
1648 event_kind::PRIVATE_DIRECT_MESSAGE as i32,
1649 event_kind::FILE_ATTACHMENT as i32,
1650 chat_identifier
1651 ],
1652 |row| row.get(0),
1653 )
1654 .map_err(|e| format!("query unread_count_for_chat: {e}"))?;
1655 Ok(count as u32)
1656}
1657
1658#[derive(Debug, PartialEq)]
1662pub enum UnreadMark {
1663 NoOp,
1665 Clear,
1667 Anchor(String),
1669}
1670
1671pub async fn compute_unread_anchor(chat_identifier: &str) -> Result<UnreadMark, String> {
1675 let conn = super::get_db_connection_guard_static()?;
1676 let (k0, k1, k2) = (
1677 event_kind::CHAT_MESSAGE as i32,
1678 event_kind::PRIVATE_DIRECT_MESSAGE as i32,
1679 event_kind::FILE_ATTACHMENT as i32,
1680 );
1681 let (target_ts, newest_ts): (Option<i64>, Option<i64>) = conn
1683 .query_row(
1684 "SELECT MAX(CASE WHEN e.mine = 0 THEN e.created_at END), MAX(e.created_at) \
1685 FROM events e JOIN chats c ON e.chat_id = c.id \
1686 WHERE c.chat_identifier = ?1 AND e.kind IN (?2, ?3, ?4)",
1687 rusqlite::params![chat_identifier, k0, k1, k2],
1688 |row| Ok((row.get(0)?, row.get(1)?)),
1689 )
1690 .map_err(|e| format!("unread anchor target: {e}"))?;
1691
1692 let target_ts = match target_ts {
1693 Some(t) => t,
1694 None => return Ok(UnreadMark::NoOp), };
1696 if newest_ts.map_or(false, |n| n > target_ts) {
1697 return Ok(UnreadMark::NoOp); }
1699
1700 let anchor_id: Option<String> = conn
1701 .query_row(
1702 "SELECT e.id FROM events e JOIN chats c ON e.chat_id = c.id \
1703 WHERE c.chat_identifier = ?1 AND e.kind IN (?2, ?3, ?4) AND e.created_at < ?5 \
1704 ORDER BY e.created_at DESC LIMIT 1",
1705 rusqlite::params![chat_identifier, k0, k1, k2, target_ts],
1706 |row| row.get(0),
1707 )
1708 .optional()
1709 .map_err(|e| format!("unread anchor prev: {e}"))?;
1710
1711 Ok(match anchor_id {
1712 Some(id) => UnreadMark::Anchor(id),
1713 None => UnreadMark::Clear,
1714 })
1715}
1716
1717pub async fn flush_message_batch(
1724 chat_id: &str,
1725 pending: &mut Vec<&Message>,
1726 session: &crate::state::SessionGuard,
1727) {
1728 if pending.is_empty() {
1729 return;
1730 }
1731 if !session.is_valid() {
1732 pending.clear();
1733 return;
1734 }
1735 if let Err(e) = save_messages_batch(chat_id, pending, Some(session)).await {
1736 crate::log_warn!("[DB] batch flush failed for {}: {}", chat_id, e);
1737 }
1738 pending.clear();
1739}
1740
1741pub async fn save_chat_messages(chat_id: &str, messages: &[Message]) -> Result<(), String> {
1743 if messages.is_empty() {
1744 return Ok(());
1745 }
1746 let refs: Vec<&Message> = messages.iter().collect();
1747 save_messages_batch(chat_id, &refs, None).await.map(|_| ())
1748}
1749
1750#[cfg(test)]
1751mod tests {
1752 use super::*;
1753 use crate::stored_event::SystemEventType;
1754
1755 static TEST_COUNTER: std::sync::atomic::AtomicU32 = std::sync::atomic::AtomicU32::new(71000);
1756
1757 fn make_test_npub(n: u32) -> String {
1758 const BECH32: &[u8] = b"qpzry9x8gf2tvdw0s3jn54khce6mua7l";
1759 let mut payload = vec![b'q'; 58];
1760 let mut x = n as u64;
1761 let mut i = 58;
1762 while x > 0 && i > 0 {
1763 i -= 1;
1764 payload[i] = BECH32[(x as usize) % 32];
1765 x /= 32;
1766 }
1767 format!("npub1{}", std::str::from_utf8(&payload).unwrap())
1768 }
1769
1770 fn init_test_db() -> (tempfile::TempDir, std::sync::MutexGuard<'static, ()>) {
1771 let guard = crate::db::DB_TEST_GUARD.lock().unwrap_or_else(|e| e.into_inner());
1772 crate::db::close_database();
1773 crate::db::clear_id_caches();
1776 let tmp = tempfile::tempdir().unwrap();
1777 let n = TEST_COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1778 let account = make_test_npub(n);
1779 std::fs::create_dir_all(tmp.path().join(&account)).unwrap();
1780 crate::db::set_app_data_dir(tmp.path().to_path_buf());
1781 crate::db::set_current_account(account.clone()).unwrap();
1782 crate::db::init_database(&account).unwrap();
1783 (tmp, guard)
1784 }
1785
1786 fn now_secs() -> u64 {
1787 std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs()
1788 }
1789
1790 #[tokio::test]
1797 async fn batch_resolver_fills_quotes_for_a_synced_page() {
1798 let (_tmp, _guard) = init_test_db();
1799 let chat = "channel_reply_batch";
1800 let author = make_test_npub(90_001);
1801
1802 let mut parent = Message::default();
1804 parent.id = "parent_evt".to_string();
1805 parent.content = "the original message".to_string();
1806 parent.npub = Some(author.clone());
1807 parent.at = now_secs() - 5_000;
1808 save_message(chat, &parent).await.unwrap();
1809
1810 let mut reply = Message::default();
1812 reply.id = "reply_evt".to_string();
1813 reply.content = "replying now".to_string();
1814 reply.replied_to = "parent_evt".to_string();
1815 reply.at = now_secs();
1816 let mut plain = Message::default();
1817 plain.id = "plain_evt".to_string();
1818 plain.content = "unrelated".to_string();
1819 plain.at = now_secs();
1820 let mut orphan = Message::default();
1822 orphan.id = "orphan_evt".to_string();
1823 orphan.replied_to = "never_seen_evt".to_string();
1824 orphan.at = now_secs();
1825
1826 populate_reply_contexts(vec![&mut reply, &mut plain, &mut orphan]).await.unwrap();
1827
1828 assert_eq!(
1829 reply.replied_to_content.as_deref(),
1830 Some("the original message"),
1831 "the synced reply carries its parent's content"
1832 );
1833 assert_eq!(reply.replied_to_npub.as_deref(), Some(author.as_str()), "and its parent's author");
1834 assert!(plain.replied_to_content.is_none(), "a non-reply is untouched");
1835 assert!(orphan.replied_to_content.is_none(), "an unresolvable parent leaves the quote empty");
1836 }
1837
1838 #[tokio::test]
1841 async fn batch_resolver_is_a_noop_without_replies() {
1842 let (_tmp, _guard) = init_test_db();
1843 populate_reply_contexts(vec![]).await.unwrap();
1844 let mut plain = Message::default();
1845 plain.id = "solo".to_string();
1846 populate_reply_contexts(vec![&mut plain]).await.unwrap();
1847 assert!(plain.replied_to_content.is_none());
1848 }
1849
1850 #[tokio::test]
1854 async fn system_event_stamps_authenticated_time_clamped_to_now() {
1855 let (_tmp, _guard) = init_test_db();
1856 let chat = "channel_ch2_timestamp";
1857 let before = now_secs();
1858 let past = before - 100_000;
1859
1860 save_system_event_at("ev_past", chat, SystemEventType::MemberJoined, "npubX", None, past, None, None).await.unwrap();
1861 save_system_event_at("ev_future", chat, SystemEventType::MemberJoined, "npubX", None, before + 100_000, None, None).await.unwrap();
1862 let after = now_secs();
1863
1864 let evs = get_system_events_for_chat(chat).unwrap();
1865 let past_ev = evs.iter().find(|e| e.id == "ev_past").expect("past event saved");
1866 assert_eq!(past_ev.created_at, past, "historical join keeps its real (authenticated) timestamp");
1867
1868 let fut_ev = evs.iter().find(|e| e.id == "ev_future").expect("future event saved");
1869 assert!(fut_ev.created_at >= before && fut_ev.created_at <= after,
1870 "future-dated event clamped to local now ({} not in {}..={})", fut_ev.created_at, before, after);
1871 }
1872
1873 #[tokio::test]
1876 async fn event_delete_context_resolves_from_db() {
1877 let (_tmp, _guard) = init_test_db();
1878 let chat = "npub1contactdc";
1879
1880 let mine_msg = Message { id: "dc_mine".into(), content: "x".into(), at: 1_000, mine: true, ..Default::default() };
1881 let theirs = Message {
1882 id: "dc_theirs".into(), content: "y".into(), at: 2_000, mine: false,
1883 npub: Some("npub1sender".to_string()),
1884 ..Default::default()
1885 };
1886 save_message(chat, &mine_msg).await.unwrap();
1887 save_message(chat, &theirs).await.unwrap();
1888
1889 let (chat_id, mine, _author) = event_delete_context("dc_mine").unwrap().expect("own row resolves");
1890 assert_eq!(chat_id, chat);
1891 assert!(mine);
1892
1893 let (chat_id, mine, author) = event_delete_context("dc_theirs").unwrap().expect("contact row resolves");
1894 assert_eq!(chat_id, chat);
1895 assert!(!mine);
1896 assert_eq!(author.as_deref(), Some("npub1sender"));
1897
1898 assert!(event_delete_context("dc_absent").unwrap().is_none(), "unknown id is None, not an error");
1899 }
1900
1901 #[tokio::test]
1904 async fn unread_counts_match_walk_back_semantics() {
1905 let (_tmp, _guard) = init_test_db();
1906 let chat = "npub1contactdm";
1907 let mk = |id: &str, secs: u64, mine: bool| Message {
1909 id: id.into(), content: "x".into(), at: secs * 1000, mine,
1910 npub: (!mine).then(|| "npub1sender".to_string()),
1911 ..Default::default()
1912 };
1913 let unread = || async { unread_counts().await.unwrap().get(chat).copied().unwrap_or(0) };
1914
1915 for i in 0..6u64 {
1917 save_message(chat, &mk(&format!("m{i}"), 1000 + i, false)).await.unwrap();
1918 }
1919 assert_eq!(unread().await, 6, "never-read backlog counts all 6");
1920
1921 save_message(chat, &mk("m6", 2000, false)).await.unwrap();
1923 save_message(chat, &mk("m7", 2001, false)).await.unwrap();
1924 assert_eq!(unread().await, 8, "6 backlog + 2 new = 8");
1925
1926 save_message(chat, &mk("mine", 2002, true)).await.unwrap();
1928 assert_eq!(unread().await, 0, "own message = read up to here");
1929
1930 save_message(chat, &mk("m8", 2003, false)).await.unwrap();
1932 assert_eq!(unread().await, 1, "one new after our send");
1933
1934 {
1936 let conn = crate::db::get_write_connection_guard_static().unwrap();
1937 conn.execute(
1938 "UPDATE chats SET last_read = ?1 WHERE chat_identifier = ?2",
1939 rusqlite::params!["m8", chat],
1940 ).unwrap();
1941 }
1942 assert_eq!(unread().await, 0, "last_read=m8 clears all");
1943 save_message(chat, &mk("m9", 2004, false)).await.unwrap();
1944 assert_eq!(unread().await, 1, "one arrival after last_read");
1945 }
1946
1947 #[tokio::test]
1950 async fn unread_count_for_chat_matches_the_map() {
1951 let (_tmp, _guard) = init_test_db();
1952 let chat = "npub1reconcile";
1953 let mk = |id: &str, secs: u64, mine: bool| Message {
1954 id: id.into(), content: "x".into(), at: secs * 1000, mine,
1955 npub: (!mine).then(|| "npub1sender".to_string()),
1956 ..Default::default()
1957 };
1958 let agree = || async {
1959 let map = unread_counts().await.unwrap().get(chat).copied().unwrap_or(0);
1960 let one = unread_count_for_chat(chat).await.unwrap();
1961 assert_eq!(map, one, "single-chat query diverged from the map");
1962 one
1963 };
1964
1965 for i in 0..4u64 { save_message(chat, &mk(&format!("m{i}"), 1000 + i, false)).await.unwrap(); }
1966 assert_eq!(agree().await, 4);
1967 save_message(chat, &mk("mine", 1010, true)).await.unwrap();
1968 assert_eq!(agree().await, 0);
1969 save_message(chat, &mk("after", 1011, false)).await.unwrap();
1970 assert_eq!(agree().await, 1);
1971 {
1972 let conn = crate::db::get_write_connection_guard_static().unwrap();
1973 conn.execute("UPDATE chats SET last_read = ?1 WHERE chat_identifier = ?2",
1974 rusqlite::params!["after", chat]).unwrap();
1975 }
1976 assert_eq!(agree().await, 0);
1977 assert_eq!(unread_count_for_chat("npub1nonexistent").await.unwrap(), 0);
1979 }
1980
1981 #[tokio::test]
1982 async fn attachments_table_round_trip_dedup_and_cascade() {
1983 let (_tmp, _guard) = init_test_db();
1984 let att = |id: &str, name: &str| Attachment {
1986 id: id.into(), url: format!("https://blossom/{id}"), name: name.into(),
1987 extension: "png".into(), size: 42, downloaded: false, ..Default::default()
1988 };
1989 let msg = |mid: &str, secs: u64, atts: Vec<Attachment>| Message {
1990 id: mid.into(), content: String::new(), at: secs * 1000, mine: false,
1991 npub: Some("npub1sender".into()), attachments: atts, ..Default::default()
1992 };
1993
1994 save_message("npub1att", &msg("m1", 1000, vec![att("hashA", "a.png"), att("hashB", "b.png")])).await.unwrap();
1996 let got = crate::db::attachments::get_attachments_for_event("m1").unwrap();
1997 assert_eq!(got.len(), 2);
1998 assert_eq!((got[0].id.as_str(), got[1].id.as_str()), ("hashA", "hashB"), "att_index order");
1999 assert_eq!(got[0].name, "a.png");
2000 assert_eq!(got[0].size, 42);
2001 assert!(!got[0].downloaded);
2002
2003 crate::db::attachments::set_attachment_downloaded("m1", "hashA", true, "/tmp/a.png").unwrap();
2005 let got = crate::db::attachments::get_attachments_for_event("m1").unwrap();
2006 assert!(got[0].downloaded && got[0].path == "/tmp/a.png");
2007 assert!(!got[1].downloaded, "sibling attachment untouched");
2008
2009 save_message("npub1att", &msg("m2", 1001, vec![att("hashA", "a-again.png")])).await.unwrap();
2011 let affected = crate::db::attachments::backfill_downloaded_by_hash("hashA", "/tmp/a.png", "m1").unwrap();
2012 assert_eq!(affected, vec!["m2".to_string()]);
2013 assert!(crate::db::attachments::get_attachments_for_event("m2").unwrap()[0].downloaded);
2014
2015 let chat_int = crate::db::id_cache::get_or_create_chat_id("npub1att").unwrap();
2018 let views = get_message_views(chat_int, 10, 0).await.unwrap();
2019 let m1 = views.iter().find(|m| m.id == "m1").unwrap();
2020 assert_eq!(m1.attachments.len(), 2);
2021
2022 delete_event("m1").await.unwrap();
2024 assert!(crate::db::attachments::get_attachments_for_event("m1").unwrap().is_empty(), "ON DELETE CASCADE");
2025 }
2026
2027 #[tokio::test]
2030 async fn attachment_download_state_is_monotonic_across_resaves() {
2031 let (_tmp, _guard) = init_test_db();
2032 let att = |id: &str, downloaded: bool, path: &str| Attachment {
2033 id: id.into(), url: "u".into(), name: "f.png".into(), extension: "png".into(),
2034 size: 1, downloaded, path: path.into(), ..Default::default()
2035 };
2036 let msg = |atts: Vec<Attachment>| Message {
2037 id: "dl1".into(), content: String::new(), at: 1_000_000, mine: false,
2038 npub: Some("npub1s".into()), attachments: atts, ..Default::default()
2039 };
2040
2041 save_message("npub1dl", &msg(vec![att("nonceid", false, "")])).await.unwrap();
2043 crate::db::attachments::set_attachment_downloaded("dl1", "nonceid", true, "/tmp/f.png").unwrap();
2044
2045 save_message("npub1dl", &msg(vec![att("nonceid", false, "")])).await.unwrap();
2047 let got = crate::db::attachments::get_attachments_for_event("dl1").unwrap();
2048 assert!(got[0].downloaded && got[0].path == "/tmp/f.png", "re-delivery preserves the download");
2049
2050 save_message("npub1dl", &msg(vec![att("contenthash", true, "/tmp/f.png")])).await.unwrap();
2052 let got = crate::db::attachments::get_attachments_for_event("dl1").unwrap();
2053 assert_eq!(got[0].id, "contenthash", "hash rewritten nonce→content");
2054 assert!(got[0].downloaded && got[0].path == "/tmp/f.png");
2055
2056 save_message("npub1dl", &msg(vec![att("nonceid", false, "")])).await.unwrap();
2059 let got = crate::db::attachments::get_attachments_for_event("dl1").unwrap();
2060 assert_eq!(got[0].id, "contenthash", "content-hash key survives a later re-delivery");
2061 assert!(got[0].downloaded && got[0].path == "/tmp/f.png");
2062 }
2063
2064 #[tokio::test]
2067 async fn attachments_fall_back_to_legacy_tag_when_table_empty() {
2068 let (_tmp, _guard) = init_test_db();
2069 let a = Attachment {
2070 id: "tagonly".into(), url: "u".into(), name: "old.png".into(), extension: "png".into(),
2071 size: 7, downloaded: false, ..Default::default()
2072 };
2073 save_message("npub1old", &Message {
2074 id: "old1".into(), content: String::new(), at: 2_000_000, mine: false,
2075 npub: Some("npub1s".into()), attachments: vec![a.clone()], ..Default::default()
2076 }).await.unwrap();
2077
2078 {
2080 let conn = crate::db::get_write_connection_guard_static().unwrap();
2081 conn.execute("DELETE FROM attachments WHERE event_id='old1'", []).unwrap();
2082 let inner = serde_json::to_string(&vec![a]).unwrap();
2083 let tags = serde_json::to_string(&vec![vec!["attachments".to_string(), inner]]).unwrap();
2084 conn.execute("UPDATE events SET tags=?1 WHERE id='old1'", rusqlite::params![tags]).unwrap();
2085 }
2086 assert!(crate::db::attachments::get_attachments_for_event("old1").unwrap().is_empty(), "table empty");
2087
2088 let chat_int = crate::db::id_cache::get_or_create_chat_id("npub1old").unwrap();
2089 let views = get_message_views(chat_int, 10, 0).await.unwrap();
2090 let old = views.iter().find(|m| m.id == "old1").unwrap();
2091 assert_eq!(old.attachments.len(), 1, "attachment served from the legacy-tag fallback");
2092 assert_eq!(old.attachments[0].name, "old.png");
2093 }
2094
2095 #[tokio::test]
2098 async fn attachment_tag_strip_is_gated_on_backfill() {
2099 let (_tmp, _guard) = init_test_db();
2100 let a = Attachment { id: "h1".into(), extension: "png".into(), name: "f.png".into(), downloaded: false, ..Default::default() };
2101 let mk = |id: &str, secs: u64| Message {
2102 id: id.into(), content: String::new(), at: secs * 1000, mine: false,
2103 npub: Some("npub1s".into()), attachments: vec![a.clone()], ..Default::default()
2104 };
2105 save_message("npub1s", &mk("bf", 1000)).await.unwrap();
2106 save_message("npub1s", &mk("unbf", 2000)).await.unwrap();
2107
2108 {
2109 let conn = crate::db::get_write_connection_guard_static().unwrap();
2110 let inner = serde_json::to_string(&vec![a.clone()]).unwrap();
2111 let with_tag = |ms: &str| serde_json::to_string(&vec![
2112 vec!["ms".to_string(), ms.to_string()],
2113 vec!["attachments".to_string(), inner.clone()],
2114 ]).unwrap();
2115 conn.execute("UPDATE events SET tags=?1 WHERE id='bf'", rusqlite::params![with_tag("5")]).unwrap();
2117 conn.execute("UPDATE events SET tags=?1 WHERE id='unbf'", rusqlite::params![with_tag("6")]).unwrap();
2118 conn.execute("DELETE FROM attachments WHERE event_id='unbf'", []).unwrap();
2119
2120 let events: Vec<(String, String)> = {
2122 let mut stmt = conn.prepare("SELECT id, tags FROM events WHERE tags LIKE '%attachments%' AND id IN (SELECT DISTINCT event_id FROM attachments)").unwrap();
2123 let m = stmt.query_map([], |r| Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?))).unwrap();
2124 m.flatten().collect()
2125 };
2126 for (id, tj) in events {
2127 let mut tags: Vec<Vec<String>> = serde_json::from_str(&tj).unwrap();
2128 tags.retain(|t| t.first().map(|s| s.as_str()) != Some("attachments"));
2129 conn.execute("UPDATE events SET tags=?1 WHERE id=?2",
2130 rusqlite::params![serde_json::to_string(&tags).unwrap(), id]).unwrap();
2131 }
2132
2133 let bf: String = conn.query_row("SELECT tags FROM events WHERE id='bf'", [], |r| r.get(0)).unwrap();
2134 assert!(!bf.contains("attachments"), "backfilled event's attachments tag stripped");
2135 assert!(bf.contains("\"ms\""), "sibling ms tag survives the strip");
2136 let unbf: String = conn.query_row("SELECT tags FROM events WHERE id='unbf'", [], |r| r.get(0)).unwrap();
2137 assert!(unbf.contains("attachments"), "un-backfilled event keeps its tag (no table row)");
2138 }
2139
2140 let chat_int = crate::db::id_cache::get_or_create_chat_id("npub1s").unwrap();
2142 let views = get_message_views(chat_int, 10, 0).await.unwrap();
2143 assert_eq!(views.iter().find(|m| m.id == "unbf").unwrap().attachments.len(), 1, "fallback still renders unbf");
2144 }
2145
2146 #[tokio::test]
2150 async fn compute_unread_anchor_covers_the_cases() {
2151 let (_tmp, _guard) = init_test_db();
2152 let mk = |id: &str, secs: u64, mine: bool| Message {
2153 id: id.into(), content: "x".into(), at: secs * 1000, mine,
2154 npub: (!mine).then(|| "npub1sender".to_string()),
2155 ..Default::default()
2156 };
2157 let unread = |chat: &'static str| async move {
2158 unread_counts().await.unwrap().get(chat).copied().unwrap_or(0)
2159 };
2160 let set_lr = |chat: &str, lr: &str| {
2161 let conn = crate::db::get_write_connection_guard_static().unwrap();
2162 conn.execute("UPDATE chats SET last_read = ?1 WHERE chat_identifier = ?2",
2163 rusqlite::params![lr, chat]).unwrap();
2164 };
2165
2166 let a = "npub1anchorA";
2169 save_message(a, &mk("a_mine", 1000, true)).await.unwrap();
2170 for i in 0..8u64 { save_message(a, &mk(&format!("a{i}"), 2000 + i, false)).await.unwrap(); }
2171 assert_eq!(compute_unread_anchor(a).await.unwrap(), UnreadMark::Anchor("a6".into()));
2172 set_lr(a, "a6");
2173 assert_eq!(unread(a).await, 1, "A: newest contact message is the sole unread");
2174
2175 let b = "npub1anchorB";
2177 save_message(b, &mk("b0", 2000, false)).await.unwrap();
2178 save_message(b, &mk("b_mine", 2001, true)).await.unwrap();
2179 assert_eq!(compute_unread_anchor(b).await.unwrap(), UnreadMark::NoOp);
2180
2181 let c = "npub1anchorC";
2184 save_message(c, &mk("c0", 3000, false)).await.unwrap();
2185 save_message(c, &mk("c1", 3005, false)).await.unwrap();
2186 save_message(c, &mk("c2", 3005, false)).await.unwrap();
2187 assert_eq!(compute_unread_anchor(c).await.unwrap(), UnreadMark::Anchor("c0".into()));
2188 set_lr(c, "c0");
2189 assert_eq!(unread(c).await, 2, "C: same-second tail both count");
2190
2191 let d = "npub1anchorD";
2193 save_message(d, &mk("d0", 4000, false)).await.unwrap();
2194 assert_eq!(compute_unread_anchor(d).await.unwrap(), UnreadMark::Clear);
2195 set_lr(d, "");
2196 assert_eq!(unread(d).await, 1, "D: lone contact message surfaces");
2197
2198 let e = "npub1anchorE";
2200 save_message(e, &mk("e_mine", 5000, true)).await.unwrap();
2201 assert_eq!(compute_unread_anchor(e).await.unwrap(), UnreadMark::NoOp);
2202 }
2203
2204 #[tokio::test]
2208 async fn unread_clears_when_last_read_is_a_system_event() {
2209 let (_tmp, _guard) = init_test_db();
2210 let chat = "npub1sysevtdm";
2211 let mk = |id: &str, secs: u64| Message {
2212 id: id.into(), content: "x".into(), at: secs * 1000, mine: false,
2213 npub: Some("npub1sender".to_string()), ..Default::default()
2214 };
2215 let unread = || async { unread_counts().await.unwrap().get(chat).copied().unwrap_or(0) };
2216
2217 for i in 0..5u64 {
2220 save_message(chat, &mk(&format!("m{i}"), 1000 + i)).await.unwrap();
2221 }
2222 assert_eq!(unread().await, 5, "never-read backlog");
2223
2224 save_system_event_at("sysev", chat, SystemEventType::MemberJoined, "npubX", None, 2000, None, None).await.unwrap();
2226
2227 {
2229 let conn = crate::db::get_write_connection_guard_static().unwrap();
2230 conn.execute(
2231 "UPDATE chats SET last_read = ?1 WHERE chat_identifier = ?2",
2232 rusqlite::params!["sysev", chat],
2233 ).unwrap();
2234 }
2235 assert_eq!(unread().await, 0, "read marker on a system event still clears the badge");
2236 }
2237
2238 #[tokio::test]
2242 async fn deleting_a_message_adjusts_unread_without_wedging() {
2243 let (_tmp, _guard) = init_test_db();
2244 let chat = "npub1delunread";
2245 let mk = |id: &str, secs: u64| Message {
2246 id: id.into(), content: "x".into(), at: secs * 1000, mine: false,
2247 npub: Some("npub1sender".to_string()), ..Default::default()
2248 };
2249 let unread = || async { unread_counts().await.unwrap().get(chat).copied().unwrap_or(0) };
2250 let set_marker = |id: &str| {
2251 let conn = crate::db::get_write_connection_guard_static().unwrap();
2252 conn.execute("UPDATE chats SET last_read = ?1 WHERE chat_identifier = ?2",
2253 rusqlite::params![id, chat]).unwrap();
2254 };
2255 let marker = || -> String {
2256 let conn = crate::db::get_db_connection_guard_static().unwrap();
2257 conn.query_row("SELECT last_read FROM chats WHERE chat_identifier = ?1",
2258 rusqlite::params![chat], |r| r.get::<_, String>(0)).unwrap()
2259 };
2260
2261 for i in 0..6u64 { save_message(chat, &mk(&format!("m{i}"), 1000 + i)).await.unwrap(); }
2263 set_marker("m1");
2264 assert_eq!(unread().await, 4, "m2..m5 unread");
2265
2266 delete_event("m3").await.unwrap();
2268 assert_eq!(unread().await, 3, "one unread deleted → badge minus one");
2269 assert_eq!(marker(), "m1", "deleting an unread message leaves the marker alone");
2270
2271 delete_event("m1").await.unwrap();
2273 assert_eq!(marker(), "m0", "marker retreats to the newest survivor before it");
2274 assert_eq!(unread().await, 3, "retreat keeps the count, no collapse to 99+");
2275
2276 delete_event("m0").await.unwrap();
2278 assert_eq!(marker(), "", "no earlier survivor → marker clears");
2279 assert_eq!(unread().await, 3, "cleared marker counts only the true unread survivors");
2280 }
2281
2282 #[tokio::test]
2286 async fn edit_event_folds_into_history_on_reload() {
2287 let (_tmp, _guard) = init_test_db();
2288 let chat = "channel_edit_fold";
2289 save_message(chat, &Message {
2290 id: "orig1".into(), content: "original".into(), at: 5_000_000,
2291 npub: Some("npub1author".into()), ..Default::default()
2292 }).await.unwrap();
2293
2294 let cid = crate::db::id_cache::get_chat_id_by_identifier(chat).unwrap();
2295 let emoji = vec![crate::types::EmojiTag { shortcode: "wave".into(), url: "u/wave".into() }];
2296 save_edit_event("edit1", "orig1", "edited :wave:", &emoji, cid, None, "npub1author").await.unwrap();
2297
2298 let m = get_message_views(cid, 50, 0).await.unwrap()
2299 .into_iter().find(|m| m.id == "orig1").expect("message reloaded");
2300 assert!(m.edited, "folded edit sets the edited flag");
2301 let h = m.edit_history.as_ref().expect("history reconstructed from the edit event");
2302 assert_eq!(h.len(), 2, "original + one edit");
2303 assert_eq!(h[0].content, "original");
2304 assert_eq!(h[1].content, "edited :wave:");
2305 assert_eq!(m.content, "edited :wave:", "latest revision is the displayed content");
2306 assert_eq!(m.emoji_tags.len(), 1, "the edit's own emoji folds onto the message");
2307 assert_eq!(m.emoji_tags[0].shortcode, "wave");
2308 }
2309
2310 #[tokio::test]
2314 async fn batched_save_matches_single_save_shape() {
2315 let (_tmp, _guard) = init_test_db();
2316 let chat = "channel_batch1";
2317 let att = crate::types::Attachment {
2318 id: "atthash1".into(), extension: "png".into(), name: "a.png".into(),
2319 url: "https://x/att".into(), downloaded: false, ..Default::default()
2320 };
2321 let reaction = Reaction {
2322 id: "react_b1".into(), reference_id: "b1".into(),
2323 author_id: "npub1reactor".into(), emoji: "👍".into(), emoji_url: None,
2324 };
2325 let msgs: Vec<Message> = (0..5u64).map(|i| Message {
2327 id: format!("b{i}"), content: format!("c{i}"), at: 7_000_000,
2328 npub: Some("npub1sender".into()),
2329 attachments: if i == 2 { vec![att.clone()] } else { Vec::new() },
2330 reactions: if i == 1 { vec![reaction.clone()] } else { Vec::new() },
2331 ..Default::default()
2332 }).collect();
2333 let refs: Vec<&Message> = msgs.iter().collect();
2334
2335 let saved = save_messages_batch(chat, &refs, None).await.unwrap();
2336 assert_eq!(saved, 5, "every message written");
2337
2338 for i in 0..5u64 {
2339 assert!(event_exists(&format!("b{i}")).unwrap(), "b{i} row exists");
2340 }
2341 assert!(event_exists("react_b1").unwrap(), "reaction landed as its own kind-7 row");
2342 let atts = crate::db::attachments::get_attachments_for_event("b2").unwrap();
2343 assert_eq!(atts.len(), 1, "attachment row committed with its event");
2344 assert_eq!(atts[0].id, "atthash1");
2345
2346 let conn = crate::db::get_db_connection_guard_static().unwrap();
2348 let ids: Vec<String> = conn
2349 .prepare("SELECT id FROM events WHERE id IN ('b0','b1','b2','b3','b4') ORDER BY rowid")
2350 .unwrap()
2351 .query_map([], |r| r.get(0)).unwrap()
2352 .flatten().collect();
2353 assert_eq!(ids, vec!["b0", "b1", "b2", "b3", "b4"], "insert order preserves the rowid tiebreak");
2354 }
2355
2356 #[tokio::test]
2359 async fn batched_resave_preserves_wrapper_and_dedups_reactions() {
2360 let (_tmp, _guard) = init_test_db();
2361 let chat = "channel_batch2";
2362 let reaction = Reaction {
2363 id: "react_rs".into(), reference_id: "rs1".into(),
2364 author_id: "npub1reactor".into(), emoji: "🔥".into(), emoji_url: None,
2365 };
2366 let mut msg = Message {
2367 id: "rs1".into(), content: "hello".into(), at: 8_000_000,
2368 npub: Some("npub1sender".into()),
2369 wrapper_event_id: Some("wrap_original".into()),
2370 reactions: vec![reaction],
2371 ..Default::default()
2372 };
2373 save_message(chat, &msg).await.unwrap();
2374
2375 msg.wrapper_event_id = None;
2377 let saved = save_messages_batch(chat, &[&msg], None).await.unwrap();
2378 assert_eq!(saved, 1);
2379
2380 let conn = crate::db::get_db_connection_guard_static().unwrap();
2381 let wrapper: Option<String> = conn.query_row(
2382 "SELECT wrapper_event_id FROM events WHERE id = 'rs1'", [], |r| r.get(0),
2383 ).unwrap();
2384 assert_eq!(wrapper.as_deref(), Some("wrap_original"), "COALESCE keeps the known wrapper");
2385 let reaction_rows: i64 = conn.query_row(
2386 "SELECT COUNT(*) FROM events WHERE id = 'react_rs'", [], |r| r.get(0),
2387 ).unwrap();
2388 assert_eq!(reaction_rows, 1, "reaction row not duplicated by the re-save");
2389 }
2390
2391 #[tokio::test]
2395 async fn batching_persist_flushes_multi_chat_and_drops_on_stale_session() {
2396 let (_tmp, _guard) = init_test_db();
2397 let handler = crate::event_handler::NoOpEventHandler;
2398 let batcher = crate::event_handler::BatchingPersist::new(&handler);
2399
2400 let mk = |id: &str, npub: &str| Message {
2401 id: id.into(), content: "x".into(), at: 9_000_000,
2402 npub: Some(npub.into()), ..Default::default()
2403 };
2404 let seed = |chat: &str, m: &Message| {
2407 let m = m.clone();
2408 let chat = chat.to_string();
2409 async move {
2410 let mut st = crate::state::STATE.lock().await;
2411 st.add_message_to_participant(&chat, &m);
2412 }
2413 };
2414 let wrap_a1 = ([0xA1u8; 32], 111u64);
2415 let wrap_b1 = ([0xB1u8; 32], 222u64);
2416 let a1 = mk("bp_a1", "npub1chata");
2417 let b1 = mk("bp_b1", "npub1chatb");
2418 let a2 = mk("bp_a2", "npub1chata");
2419 seed("npub1chata", &a1).await;
2420 seed("npub1chatb", &b1).await;
2421 seed("npub1chata", &a2).await;
2422
2423 use crate::event_handler::InboundEventHandler;
2424 assert!(batcher.buffer_persist("npub1chata", &a1, Some(wrap_a1)), "batcher owns the persist");
2425 assert!(batcher.buffer_persist("npub1chatb", &b1, Some(wrap_b1)));
2426 assert!(batcher.buffer_persist("npub1chata", &a2, None));
2427 assert_eq!(batcher.buffered(), 3);
2428
2429 let ledgered = |bytes: [u8; 32]| {
2430 let id = nostr_sdk::prelude::EventId::from_byte_array(bytes);
2431 crate::db::wrappers::load_negentropy_items().unwrap().iter().any(|(e, _)| *e == id)
2432 };
2433 assert!(!ledgered(wrap_a1.0), "wrapper unledgered while its message sits buffered");
2434
2435 let session = crate::state::SessionGuard::capture();
2436 assert_eq!(batcher.flush(&session).await, 3, "all buffered messages written");
2437 assert_eq!(batcher.buffered(), 0);
2438 assert!(event_exists("bp_a1").unwrap() && event_exists("bp_b1").unwrap() && event_exists("bp_a2").unwrap());
2439 assert!(ledgered(wrap_a1.0) && ledgered(wrap_b1.0), "wrappers ledgered with the flush");
2440 let a = crate::db::id_cache::get_chat_id_by_identifier("npub1chata").unwrap();
2441 let b = crate::db::id_cache::get_chat_id_by_identifier("npub1chatb").unwrap();
2442 assert_ne!(a, b, "rows grouped under their own chats");
2443
2444 let stale = mk("bp_stale", "npub1chata");
2447 seed("npub1chata", &stale).await;
2448 let wrap_stale = ([0x5Eu8; 32], 333u64);
2449 batcher.buffer_persist("npub1chata", &stale, Some(wrap_stale));
2450 crate::state::bump_session_generation();
2451 assert_eq!(batcher.flush(&session).await, 0, "stale flush writes nothing");
2452 assert!(!event_exists("bp_stale").unwrap(), "stale message never reached the DB");
2453 assert!(!ledgered(wrap_stale.0), "dropped message's wrapper NOT ledgered — negentropy will re-deliver it");
2454 assert_eq!(batcher.buffered(), 0, "stale buffer drained, not retried into the next account");
2455 }
2456
2457 #[tokio::test]
2464 async fn buffered_message_deleted_before_flush_never_persists() {
2465 let (_tmp, _guard) = init_test_db();
2466 let handler = crate::event_handler::NoOpEventHandler;
2467 let batcher = crate::event_handler::BatchingPersist::new(&handler);
2468 use crate::event_handler::InboundEventHandler;
2469
2470 let chat = "npub1delchat";
2471 let mk = |id: &str| Message {
2472 id: id.into(), content: "x".into(), at: 9_500_000,
2473 npub: Some(chat.into()), ..Default::default()
2474 };
2475 let ledgered = |bytes: [u8; 32]| {
2476 let id = nostr_sdk::prelude::EventId::from_byte_array(bytes);
2477 crate::db::wrappers::load_negentropy_items().unwrap().iter().any(|(e, _)| *e == id)
2478 };
2479
2480 let m1 = mk("del_sametask");
2483 {
2484 let mut st = crate::state::STATE.lock().await;
2485 st.add_message_to_participant(chat, &m1);
2486 }
2487 batcher.buffer_persist(chat, &m1, Some(([0xD1u8; 32], 444)));
2488 {
2489 let mut st = crate::state::STATE.lock().await;
2490 st.remove_message("del_sametask");
2491 }
2492 batcher.on_message_deleted(chat, "del_sametask");
2493 assert_eq!(batcher.buffered(), 0, "deletion purges the buffered target");
2494
2495 let m2 = mk("del_crosstask");
2498 {
2499 let mut st = crate::state::STATE.lock().await;
2500 st.add_message_to_participant(chat, &m2);
2501 }
2502 batcher.buffer_persist(chat, &m2, Some(([0xD2u8; 32], 555)));
2503 crate::state::note_message_deleted("del_crosstask");
2504
2505 let m3 = mk("evicted_ok");
2507 let wrap_evicted = ([0xE0u8; 32], 666u64);
2508 {
2509 let mut st = crate::state::STATE.lock().await;
2510 st.add_message_to_participant(chat, &m3);
2511 st.remove_message("evicted_ok");
2512 }
2513 batcher.buffer_persist(chat, &m3, Some(wrap_evicted));
2514
2515 let session = crate::state::SessionGuard::capture();
2516 assert_eq!(batcher.flush(&session).await, 1, "tombstoned target dropped, evicted message written");
2517 assert!(!event_exists("del_sametask").unwrap(), "purged message never persisted");
2518 assert!(!event_exists("del_crosstask").unwrap(), "tombstoned message never persisted");
2519 assert!(event_exists("evicted_ok").unwrap(), "evicted-but-not-deleted message persisted");
2520 assert!(ledgered(wrap_evicted.0), "evicted message's wrapper ledgered with it");
2521 }
2522}