1use std::path::Path;
2use std::sync::{Arc, Mutex};
3
4use nostr::{Event, Filter, Keys, PublicKey, UnsignedEvent};
5use nostr_double_ratchet::Invite;
6use rusqlite::{params, Connection, OptionalExtension};
7
8#[cfg(test)]
9use crate::AppKeys;
10use crate::{
11 invite_unsigned_event, is_app_keys_event, parse_invite_url, resolve_invite_owner,
12 ProtocolAcceptInviteBlock, ProtocolAcceptInviteOutcome, ProtocolDecryptedMessage,
13 ProtocolEffect, ProtocolEngine, ProtocolRetryBatch, SharedConnection, SqliteStorageAdapter,
14 UnixSeconds, APP_KEYS_EVENT_KIND, CHAT_MESSAGE_KIND, INVITE_EVENT_KIND, INVITE_RESPONSE_KIND,
15 MESSAGE_EVENT_KIND,
16};
17
18const SCHEMA: &str = r#"
19CREATE TABLE IF NOT EXISTS private_chat_threads (
20 chat_id TEXT PRIMARY KEY,
21 display_name TEXT NOT NULL,
22 avatar_seed TEXT NOT NULL,
23 updated_at_secs INTEGER NOT NULL DEFAULT 0
24);
25
26CREATE TABLE IF NOT EXISTS private_chat_messages (
27 chat_id TEXT NOT NULL,
28 id TEXT NOT NULL,
29 body TEXT NOT NULL,
30 is_outgoing INTEGER NOT NULL,
31 created_at_secs INTEGER NOT NULL,
32 delivery TEXT NOT NULL,
33 source_event_id TEXT,
34 PRIMARY KEY (chat_id, id)
35);
36
37CREATE INDEX IF NOT EXISTS private_chat_recent_idx
38 ON private_chat_messages(chat_id, created_at_secs, id);
39
40CREATE UNIQUE INDEX IF NOT EXISTS private_chat_source_event_idx
41 ON private_chat_messages(source_event_id)
42 WHERE source_event_id IS NOT NULL;
43
44CREATE TABLE IF NOT EXISTS private_chat_seen_events (
45 event_id TEXT PRIMARY KEY
46);
47
48CREATE TABLE IF NOT EXISTS ndr_kv (
49 owner_pubkey_hex TEXT NOT NULL,
50 device_pubkey_hex TEXT NOT NULL,
51 key TEXT NOT NULL,
52 value TEXT NOT NULL,
53 PRIMARY KEY (owner_pubkey_hex, device_pubkey_hex, key)
54);
55"#;
56
57#[derive(Clone, Debug, PartialEq, Eq)]
58pub enum DirectMessageDelivery {
59 Pending,
60 Sent,
61 Received,
62 Failed,
63}
64
65impl DirectMessageDelivery {
66 fn as_str(&self) -> &'static str {
67 match self {
68 Self::Pending => "pending",
69 Self::Sent => "sent",
70 Self::Received => "received",
71 Self::Failed => "failed",
72 }
73 }
74
75 fn from_str(value: &str) -> Self {
76 match value {
77 "sent" => Self::Sent,
78 "received" => Self::Received,
79 "failed" => Self::Failed,
80 _ => Self::Pending,
81 }
82 }
83}
84
85#[derive(Clone, Debug, PartialEq, Eq)]
86pub struct DirectMessageSnapshot {
87 pub id: String,
88 pub chat_id: String,
89 pub body: String,
90 pub is_outgoing: bool,
91 pub created_at_secs: u64,
92 pub delivery: DirectMessageDelivery,
93}
94
95#[derive(Clone, Debug, PartialEq, Eq)]
96pub struct DirectChatSnapshot {
97 pub chat_id: String,
98 pub last_message_preview: String,
99 pub last_message_at: u64,
100 pub unread_count: u32,
101}
102
103#[derive(Clone, Debug, PartialEq, Eq)]
104pub struct DirectThreadSnapshot {
105 pub chat: DirectChatSnapshot,
106 pub messages: Vec<DirectMessageSnapshot>,
107}
108
109#[derive(Clone, Debug)]
110pub enum DirectMessageCommand {
111 Publish(Event),
112 Subscribe {
113 subscription_id: String,
114 filters: Vec<Filter>,
115 durable: bool,
116 },
117}
118
119#[derive(Clone, Debug)]
120pub enum DirectInviteAcceptanceOutcome {
121 Accepted {
122 thread: DirectThreadSnapshot,
123 commands: Vec<DirectMessageCommand>,
124 },
125 PendingOwnerRoster {
126 owner_pubkey: String,
127 device_pubkey: String,
128 commands: Vec<DirectMessageCommand>,
129 },
130}
131
132pub struct DirectMessageService {
133 conn: SharedConnection,
134 protocol_engine: Option<ProtocolEngine>,
135 owner_public_key: Option<PublicKey>,
136 relay_subscription_key: Option<String>,
137 last_error: Option<String>,
138}
139
140impl DirectMessageService {
141 pub fn memory() -> Self {
142 let service = Self {
143 conn: Arc::new(Mutex::new(Connection::open_in_memory().unwrap())),
144 protocol_engine: None,
145 owner_public_key: None,
146 relay_subscription_key: None,
147 last_error: None,
148 };
149 service.ensure_schema();
150 service
151 }
152
153 pub fn memory_for_local_device(owner_public_key: PublicKey, device_keys: &Keys) -> Self {
154 Self::memory().with_protocol_engine_for_local_device(owner_public_key, device_keys)
155 }
156
157 pub fn open(data_dir: &Path, owner_keys: Option<&Keys>) -> Self {
158 match owner_keys {
159 Some(keys) => Self::open_for_local_device(data_dir, keys.public_key(), keys),
160 None => Self::open_without_protocol_engine(data_dir),
161 }
162 }
163
164 pub fn open_for_local_device(
165 data_dir: &Path,
166 owner_public_key: PublicKey,
167 device_keys: &Keys,
168 ) -> Self {
169 Self::open_without_protocol_engine(data_dir)
170 .with_protocol_engine_for_local_device(owner_public_key, device_keys)
171 }
172
173 fn open_without_protocol_engine(data_dir: &Path) -> Self {
174 let path = data_dir.join("private-chat.sqlite3");
175 let conn = Connection::open(path).or_else(|_| Connection::open_in_memory());
176 let conn = match conn {
177 Ok(conn) => conn,
178 Err(error) => {
179 return Self {
180 conn: Arc::new(Mutex::new(Connection::open_in_memory().unwrap())),
181 protocol_engine: None,
182 owner_public_key: None,
183 relay_subscription_key: None,
184 last_error: Some(format!("Direct message store open failed: {error}")),
185 };
186 }
187 };
188 let service = Self {
189 conn: Arc::new(Mutex::new(conn)),
190 protocol_engine: None,
191 owner_public_key: None,
192 relay_subscription_key: None,
193 last_error: None,
194 };
195 service.ensure_schema();
196 service
197 }
198
199 pub fn activate(&mut self, keys: &Keys) -> Vec<DirectMessageCommand> {
200 let next = Self {
201 conn: Arc::clone(&self.conn),
202 protocol_engine: None,
203 owner_public_key: None,
204 relay_subscription_key: self.relay_subscription_key.clone(),
205 last_error: self.last_error.clone(),
206 }
207 .with_protocol_engine(keys);
208 self.protocol_engine = next.protocol_engine;
209 self.owner_public_key = next.owner_public_key;
210 self.protocol_subscription_commands()
211 }
212
213 pub fn last_error(&self) -> Option<String> {
214 self.last_error.clone()
215 }
216
217 pub fn chats(&self) -> Vec<DirectChatSnapshot> {
218 let Ok(conn) = self.conn.lock() else {
219 return Vec::new();
220 };
221 let mut stmt = match conn.prepare(
222 "SELECT t.chat_id,
223 COALESCE(m.body, ''), COALESCE(m.created_at_secs, t.updated_at_secs)
224 FROM private_chat_threads t
225 LEFT JOIN private_chat_messages m
226 ON m.chat_id = t.chat_id
227 AND m.id = (
228 SELECT id
229 FROM private_chat_messages
230 WHERE chat_id = t.chat_id
231 ORDER BY created_at_secs DESC, id DESC
232 LIMIT 1
233 )
234 ORDER BY COALESCE(m.created_at_secs, t.updated_at_secs) DESC, t.chat_id ASC",
235 ) {
236 Ok(stmt) => stmt,
237 Err(_) => return Vec::new(),
238 };
239 let rows = match stmt.query_map([], |row| {
240 Ok(DirectChatSnapshot {
241 chat_id: row.get(0)?,
242 last_message_preview: row.get(1)?,
243 last_message_at: row.get::<_, i64>(2)?.max(0) as u64,
244 unread_count: 0,
245 })
246 }) {
247 Ok(rows) => rows,
248 Err(_) => return Vec::new(),
249 };
250 rows.filter_map(Result::ok).collect()
251 }
252
253 pub fn thread(&self, chat_id: &str) -> Option<DirectThreadSnapshot> {
254 let chat_id = normalize_pubkey(chat_id).ok()?;
255 let chat = self
256 .chats()
257 .into_iter()
258 .find(|chat| chat.chat_id == chat_id)
259 .unwrap_or_else(|| chat_snapshot_for_pubkey(&chat_id));
260 let messages = self.messages(&chat_id, 160);
261 Some(DirectThreadSnapshot { chat, messages })
262 }
263
264 pub fn open_chat(
265 &mut self,
266 peer_input: &str,
267 keys: &Keys,
268 ) -> Result<(DirectThreadSnapshot, Vec<DirectMessageCommand>), String> {
269 let public_key = match PublicKey::parse(peer_input) {
270 Ok(public_key) => public_key,
271 Err(_) => return self.accept_invite(peer_input, keys),
272 };
273 let chat_id = public_key.to_hex();
274 self.ensure_thread(&chat_id, unix_now());
275 let commands = self.protocol_subscription_commands();
276 let thread = self
277 .thread(&chat_id)
278 .ok_or_else(|| "Chat open failed".to_string())?;
279 Ok((thread, commands))
280 }
281
282 pub fn accept_invite(
283 &mut self,
284 invite_input: &str,
285 keys: &Keys,
286 ) -> Result<(DirectThreadSnapshot, Vec<DirectMessageCommand>), String> {
287 match self.accept_invite_with_status(invite_input, keys)? {
288 DirectInviteAcceptanceOutcome::Accepted { thread, commands } => {
289 Ok((thread, commands))
290 }
291 DirectInviteAcceptanceOutcome::PendingOwnerRoster {
292 owner_pubkey,
293 device_pubkey,
294 ..
295 } => Err(format!(
296 "Invite owner device list is not available yet for owner {owner_pubkey} and device {device_pubkey}; use accept_invite_with_status, process its discovery commands, and retry"
297 )),
298 }
299 }
300
301 pub fn accept_invite_with_status(
302 &mut self,
303 invite_input: &str,
304 _keys: &Keys,
305 ) -> Result<DirectInviteAcceptanceOutcome, String> {
306 let invite = parse_direct_invite_input(invite_input)?;
307 let owner = resolve_invite_owner(&invite, None).map_err(|error| error.to_string())?;
308 let outcome = {
309 let engine = self
310 .protocol_engine
311 .as_mut()
312 .ok_or_else(|| "Direct message runtime is not ready".to_string())?;
313 engine
314 .accept_invite(&invite, Some(owner))
315 .map_err(|error| error.to_string())?
316 };
317 let result = match outcome {
318 ProtocolAcceptInviteOutcome::Accepted(result) => result,
319 ProtocolAcceptInviteOutcome::Blocked(
320 ProtocolAcceptInviteBlock::MissingOwnerRoster {
321 owner_pubkey,
322 device_pubkey,
323 },
324 ) => {
325 let commands = vec![DirectMessageCommand::Subscribe {
326 subscription_id: format!("iris-native-app-keys-{}", owner_pubkey.to_hex()),
327 filters: vec![Filter::new()
328 .kind(nostr::Kind::from(APP_KEYS_EVENT_KIND as u16))
329 .author(owner_pubkey)
330 .limit(16)],
331 durable: true,
332 }];
333 return Ok(DirectInviteAcceptanceOutcome::PendingOwnerRoster {
334 owner_pubkey: owner_pubkey.to_hex(),
335 device_pubkey: device_pubkey.to_hex(),
336 commands,
337 });
338 }
339 ProtocolAcceptInviteOutcome::Blocked(
340 ProtocolAcceptInviteBlock::UnauthorizedDevice { .. },
341 ) => {
342 return Err("Invite device is not authorized by its claimed owner".to_string());
343 }
344 };
345 let chat_id = result.owner_pubkey.to_hex();
346 self.ensure_thread(&chat_id, unix_now());
347 let mut commands = self.commands_from_effects(result.effects);
348 commands.extend(self.protocol_subscription_commands());
349 let thread = self
350 .thread(&chat_id)
351 .ok_or_else(|| "Invite chat open failed".to_string())?;
352 Ok(DirectInviteAcceptanceOutcome::Accepted { thread, commands })
353 }
354
355 pub fn send_message(
356 &mut self,
357 chat_id: &str,
358 body: &str,
359 _keys: &Keys,
360 ) -> Result<Vec<DirectMessageCommand>, String> {
361 let body = body.trim();
362 if body.is_empty() {
363 return Ok(Vec::new());
364 }
365 let public_key = PublicKey::parse(chat_id).map_err(|error| error.to_string())?;
366 let chat_id = public_key.to_hex();
367 self.ensure_thread(&chat_id, unix_now());
368 let engine = self
369 .protocol_engine
370 .as_mut()
371 .ok_or_else(|| "Direct message runtime is not ready".to_string())?;
372 let result = engine
373 .send_direct_text(public_key, &chat_id, body, None, UnixSeconds(unix_now()))
374 .map_err(|error| error.to_string())?;
375 let delivery = if result.event_ids.is_empty() {
376 DirectMessageDelivery::Pending
377 } else {
378 DirectMessageDelivery::Sent
379 };
380 self.insert_message(
381 &chat_id,
382 &result.message_id,
383 body,
384 true,
385 unix_now(),
386 delivery,
387 None,
388 );
389 Ok(self.commands_from_effects(result.effects))
390 }
391
392 pub fn process_event(&mut self, event: Event, _keys: &Keys) -> Vec<DirectMessageCommand> {
393 let event_id = event.id.to_hex();
394 if self.seen_event(&event_id) {
395 return Vec::new();
396 }
397 let Some(engine) = self.protocol_engine.as_mut() else {
398 return Vec::new();
399 };
400 let kind = event.kind.as_u16() as u32;
401 let mut effects = Vec::new();
402 let mut retry_batch = ProtocolRetryBatch::default();
403 let mut decrypted = None;
404
405 let processed = match kind {
406 APP_KEYS_EVENT_KIND if is_app_keys_event(&event) => {
407 match engine.ingest_app_keys_event(&event) {
408 Ok(batch) => {
409 retry_batch = batch;
410 true
411 }
412 Err(error) => {
413 self.last_error =
414 Some(format!("Direct message device roster failed: {error}"));
415 false
416 }
417 }
418 }
419 INVITE_EVENT_KIND => match engine.observe_invite_event(&event) {
420 Ok(batch) => {
421 retry_batch = batch;
422 true
423 }
424 Err(_) => false,
425 },
426 INVITE_RESPONSE_KIND => match engine.observe_invite_response_event(&event) {
427 Ok(batch) => {
428 retry_batch = batch;
429 true
430 }
431 Err(_) => false,
432 },
433 MESSAGE_EVENT_KIND => match engine.process_direct_message_event(&event) {
434 Ok(message) => {
435 decrypted = message;
436 true
437 }
438 Err(_) => false,
439 },
440 _ => false,
441 };
442
443 if !processed {
444 return Vec::new();
445 }
446 let app_keys_owner =
447 (kind == APP_KEYS_EVENT_KIND && is_app_keys_event(&event)).then_some(event.pubkey);
448 self.mark_seen_event(&event_id);
449 if let Some(message) = decrypted {
450 self.apply_decrypted_protocol_message(message);
451 }
452 effects.extend(self.effects_from_retry_batch(retry_batch));
453 let mut commands = self.commands_from_effects(effects);
454 if app_keys_owner.is_some() {
455 commands.extend(self.protocol_subscription_commands());
456 }
457 commands
458 }
459
460 pub fn mobile_push_message_author_pubkeys(&self) -> Vec<String> {
461 let Some(engine) = self.protocol_engine.as_ref() else {
462 return Vec::new();
463 };
464 let mut authors = engine
465 .known_message_author_pubkeys()
466 .into_iter()
467 .map(|pubkey| pubkey.to_hex())
468 .collect::<Vec<_>>();
469 authors.sort();
470 authors.dedup();
471 authors
472 }
473
474 pub fn local_invite_event(&self, device_keys: &Keys) -> Option<Event> {
475 let invite = self.protocol_engine.as_ref()?.local_invite()?;
476 if invite.inviter_device_pubkey.to_bytes() != device_keys.public_key().to_bytes() {
477 return None;
478 }
479 invite_unsigned_event(&invite)
480 .ok()?
481 .sign_with_keys(device_keys)
482 .ok()
483 }
484
485 fn subscription_command(&mut self) -> Option<DirectMessageCommand> {
486 let engine = self.protocol_engine.as_ref()?;
487 let authors = engine
488 .known_message_author_pubkeys()
489 .into_iter()
490 .chain(self.owner_public_key)
491 .collect::<Vec<_>>();
492 let mut author_hexes = authors.iter().map(PublicKey::to_hex).collect::<Vec<_>>();
493 author_hexes.sort();
494 author_hexes.dedup();
495 let key = author_hexes.join(",");
496 if key.is_empty() || self.relay_subscription_key.as_deref() == Some(key.as_str()) {
497 return None;
498 }
499 self.relay_subscription_key = Some(key);
500
501 let public_keys = author_hexes
502 .iter()
503 .filter_map(|hex| PublicKey::parse(hex).ok())
504 .collect::<Vec<_>>();
505 let filter = Filter::new()
506 .authors(public_keys)
507 .kinds([
508 nostr::Kind::from(MESSAGE_EVENT_KIND as u16),
509 nostr::Kind::from(INVITE_EVENT_KIND as u16),
510 nostr::Kind::from(INVITE_RESPONSE_KIND as u16),
511 nostr::Kind::from(APP_KEYS_EVENT_KIND as u16),
512 ])
513 .limit(500);
514 Some(DirectMessageCommand::Subscribe {
515 subscription_id: "iris-native-private-chat".to_string(),
516 filters: vec![filter],
517 durable: true,
518 })
519 }
520
521 fn with_protocol_engine(self, keys: &Keys) -> Self {
522 self.with_protocol_engine_for_local_device(keys.public_key(), keys)
523 }
524
525 fn with_protocol_engine_for_local_device(
526 mut self,
527 owner: PublicKey,
528 device_keys: &Keys,
529 ) -> Self {
530 let owner_hex = owner.to_hex();
531 let device_hex = device_keys.public_key().to_hex();
532 let storage = Arc::new(SqliteStorageAdapter::new(
533 Arc::clone(&self.conn),
534 owner_hex.clone(),
535 device_hex,
536 ));
537 match ProtocolEngine::load_or_create_for_local_device(storage, owner, device_keys) {
538 Ok(engine) => {
539 self.protocol_engine = Some(engine);
540 self.owner_public_key = Some(owner);
541 }
542 Err(error) => self.last_error = Some(format!("Direct message init failed: {error}")),
543 }
544 self
545 }
546
547 fn protocol_subscription_commands(&mut self) -> Vec<DirectMessageCommand> {
548 self.subscription_command().into_iter().collect()
549 }
550
551 fn commands_from_effects(&mut self, effects: Vec<ProtocolEffect>) -> Vec<DirectMessageCommand> {
552 let mut commands = Vec::new();
553 for effect in effects {
554 match effect {
555 ProtocolEffect::Publish(publish) => {
556 commands.push(DirectMessageCommand::Publish(publish.event));
557 }
558 }
559 }
560 commands
561 }
562
563 fn effects_from_retry_batch(&mut self, batch: ProtocolRetryBatch) -> Vec<ProtocolEffect> {
564 let mut effects = batch.effects;
565 effects.extend(batch.group_result.effects);
566 for message in batch.direct_messages {
567 self.apply_decrypted_protocol_message(message);
568 }
569 effects
570 }
571
572 fn apply_decrypted_protocol_message(&mut self, message: ProtocolDecryptedMessage) {
573 self.apply_decrypted(
574 message.sender,
575 message.conversation_owner,
576 &message.content,
577 message.event_id,
578 );
579 }
580
581 fn apply_decrypted(
582 &mut self,
583 sender: PublicKey,
584 conversation_owner: Option<PublicKey>,
585 content: &str,
586 source_event_id: Option<String>,
587 ) {
588 let Some(rumor) = parse_runtime_rumor(content) else {
589 return;
590 };
591 if rumor.kind != CHAT_MESSAGE_KIND {
592 return;
593 }
594 let local_owner = self.owner_public_key;
595 let peer = if local_owner == Some(sender) {
596 conversation_owner.unwrap_or(sender)
597 } else {
598 sender
599 };
600 let chat_id = peer.to_hex();
601 self.ensure_thread(&chat_id, rumor.created_at_secs);
602 self.insert_message(
603 &chat_id,
604 &rumor.id,
605 &rumor.content,
606 local_owner == Some(sender),
607 rumor.created_at_secs,
608 if local_owner == Some(sender) {
609 DirectMessageDelivery::Sent
610 } else {
611 DirectMessageDelivery::Received
612 },
613 source_event_id.as_deref(),
614 );
615 }
616
617 fn ensure_schema(&self) {
618 if let Ok(conn) = self.conn.lock() {
619 let _ = conn.execute_batch(SCHEMA);
620 }
621 }
622
623 fn ensure_thread(&self, chat_id: &str, updated_at: u64) {
624 if let Ok(conn) = self.conn.lock() {
625 let _ = conn.execute(
626 "INSERT INTO private_chat_threads (chat_id, display_name, avatar_seed, updated_at_secs)
627 VALUES (?1, '', '', ?2)
628 ON CONFLICT(chat_id) DO UPDATE SET updated_at_secs = MAX(updated_at_secs, excluded.updated_at_secs)",
629 params![chat_id, updated_at as i64],
630 );
631 }
632 }
633
634 fn messages(&self, chat_id: &str, limit: usize) -> Vec<DirectMessageSnapshot> {
635 let Ok(conn) = self.conn.lock() else {
636 return Vec::new();
637 };
638 let mut stmt = match conn.prepare(
639 "SELECT id, body, is_outgoing, created_at_secs, delivery
640 FROM private_chat_messages
641 WHERE chat_id = ?1
642 ORDER BY created_at_secs DESC, id DESC
643 LIMIT ?2",
644 ) {
645 Ok(stmt) => stmt,
646 Err(_) => return Vec::new(),
647 };
648 let rows = match stmt.query_map(params![chat_id, limit as i64], |row| {
649 Ok(DirectMessageSnapshot {
650 id: row.get(0)?,
651 chat_id: chat_id.to_string(),
652 body: row.get(1)?,
653 is_outgoing: row.get::<_, i64>(2)? != 0,
654 created_at_secs: row.get::<_, i64>(3)?.max(0) as u64,
655 delivery: DirectMessageDelivery::from_str(&row.get::<_, String>(4)?),
656 })
657 }) {
658 Ok(rows) => rows,
659 Err(_) => return Vec::new(),
660 };
661 let mut messages = rows.filter_map(Result::ok).collect::<Vec<_>>();
662 messages.reverse();
663 messages
664 }
665
666 #[allow(clippy::too_many_arguments)]
667 fn insert_message(
668 &self,
669 chat_id: &str,
670 id: &str,
671 body: &str,
672 is_outgoing: bool,
673 created_at: u64,
674 delivery: DirectMessageDelivery,
675 source_event_id: Option<&str>,
676 ) {
677 if id.is_empty() {
678 return;
679 }
680 if let Ok(conn) = self.conn.lock() {
681 let _ = conn.execute(
682 "INSERT OR IGNORE INTO private_chat_messages
683 (chat_id, id, body, is_outgoing, created_at_secs, delivery, source_event_id)
684 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
685 params![
686 chat_id,
687 id,
688 body,
689 is_outgoing as i64,
690 created_at as i64,
691 delivery.as_str(),
692 source_event_id,
693 ],
694 );
695 let _ = conn.execute(
696 "UPDATE private_chat_threads SET updated_at_secs = MAX(updated_at_secs, ?2)
697 WHERE chat_id = ?1",
698 params![chat_id, created_at as i64],
699 );
700 }
701 }
702
703 fn seen_event(&self, event_id: &str) -> bool {
704 let Ok(conn) = self.conn.lock() else {
705 return true;
706 };
707 conn.query_row(
708 "SELECT 1 FROM private_chat_seen_events WHERE event_id = ?1",
709 [event_id],
710 |_| Ok(()),
711 )
712 .optional()
713 .ok()
714 .flatten()
715 .is_some()
716 }
717
718 fn mark_seen_event(&self, event_id: &str) {
719 if let Ok(conn) = self.conn.lock() {
720 let _ = conn.execute(
721 "INSERT OR IGNORE INTO private_chat_seen_events (event_id) VALUES (?1)",
722 [event_id],
723 );
724 }
725 }
726}
727
728struct RuntimeRumor {
729 id: String,
730 kind: u32,
731 content: String,
732 created_at_secs: u64,
733}
734
735fn parse_runtime_rumor(content: &str) -> Option<RuntimeRumor> {
736 let mut event = serde_json::from_str::<UnsignedEvent>(content).ok()?;
737 event.ensure_id();
738 event.verify_id().ok()?;
739 Some(RuntimeRumor {
740 id: event.id.as_ref()?.to_string(),
741 kind: event.kind.as_u16() as u32,
742 content: event.content,
743 created_at_secs: event.created_at.as_secs(),
744 })
745}
746
747fn chat_snapshot_for_pubkey(chat_id: &str) -> DirectChatSnapshot {
748 DirectChatSnapshot {
749 chat_id: chat_id.to_string(),
750 last_message_preview: String::new(),
751 last_message_at: 0,
752 unread_count: 0,
753 }
754}
755
756fn normalize_pubkey(input: &str) -> Result<String, String> {
757 PublicKey::parse(input)
758 .map(|pubkey| pubkey.to_hex())
759 .map_err(|error| error.to_string())
760}
761
762fn parse_direct_invite_input(input: &str) -> Result<Invite, String> {
763 let trimmed = input.trim();
764 if trimmed.is_empty() {
765 return Err("Invite link is required".to_string());
766 }
767 if let Ok(invite) = parse_invite_url(trimmed) {
768 return Ok(invite);
769 }
770
771 let mut candidates = vec![trimmed.to_string()];
772 if let Some((_, fragment)) = trimmed.split_once('#') {
773 candidates.push(fragment.to_string());
774 candidates.push(fragment.trim_start_matches('/').to_string());
775 candidates.extend(
776 fragment
777 .split(['/', '?', '&', '='])
778 .filter(|part| !part.trim().is_empty())
779 .map(ToString::to_string),
780 );
781 }
782 if let Some((_, query)) = trimmed.split_once('?') {
783 candidates.extend(
784 query
785 .split(['/', '?', '&', '='])
786 .filter(|part| !part.trim().is_empty())
787 .map(ToString::to_string),
788 );
789 }
790
791 for candidate in candidates {
792 let candidate = candidate.trim().trim_start_matches('/');
793 let candidate = candidate.strip_prefix("invite/").unwrap_or(candidate);
794 if candidate.is_empty() || candidate.eq_ignore_ascii_case("invite") {
795 continue;
796 }
797 for wrapped in [
798 candidate.to_string(),
799 format!("https://chat.iris.to#{candidate}"),
800 format!("https://chat.iris.to#/{candidate}"),
801 ] {
802 if let Ok(invite) = parse_invite_url(&wrapped) {
803 return Ok(invite);
804 }
805 }
806 }
807
808 parse_invite_url(trimmed).map_err(|error| error.to_string())
809}
810
811fn unix_now() -> u64 {
812 std::time::SystemTime::now()
813 .duration_since(std::time::UNIX_EPOCH)
814 .map(|duration| duration.as_secs())
815 .unwrap_or_default()
816}
817
818#[cfg(test)]
819mod tests;