Skip to main content

iris_chat_protocol/
direct_messages.rs

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;