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.created_at_secs = (
228                    SELECT MAX(created_at_secs)
229                    FROM private_chat_messages
230                    WHERE chat_id = t.chat_id
231              )
232             ORDER BY COALESCE(m.created_at_secs, t.updated_at_secs) DESC, t.chat_id ASC",
233        ) {
234            Ok(stmt) => stmt,
235            Err(_) => return Vec::new(),
236        };
237        let rows = match stmt.query_map([], |row| {
238            Ok(DirectChatSnapshot {
239                chat_id: row.get(0)?,
240                last_message_preview: row.get(1)?,
241                last_message_at: row.get::<_, i64>(2)?.max(0) as u64,
242                unread_count: 0,
243            })
244        }) {
245            Ok(rows) => rows,
246            Err(_) => return Vec::new(),
247        };
248        rows.filter_map(Result::ok).collect()
249    }
250
251    pub fn thread(&self, chat_id: &str) -> Option<DirectThreadSnapshot> {
252        let chat_id = normalize_pubkey(chat_id).ok()?;
253        let chat = self
254            .chats()
255            .into_iter()
256            .find(|chat| chat.chat_id == chat_id)
257            .unwrap_or_else(|| chat_snapshot_for_pubkey(&chat_id));
258        let messages = self.messages(&chat_id, 160);
259        Some(DirectThreadSnapshot { chat, messages })
260    }
261
262    pub fn open_chat(
263        &mut self,
264        peer_input: &str,
265        keys: &Keys,
266    ) -> Result<(DirectThreadSnapshot, Vec<DirectMessageCommand>), String> {
267        let public_key = match PublicKey::parse(peer_input) {
268            Ok(public_key) => public_key,
269            Err(_) => return self.accept_invite(peer_input, keys),
270        };
271        let chat_id = public_key.to_hex();
272        self.ensure_thread(&chat_id, unix_now());
273        let commands = self.protocol_subscription_commands();
274        let thread = self
275            .thread(&chat_id)
276            .ok_or_else(|| "Chat open failed".to_string())?;
277        Ok((thread, commands))
278    }
279
280    pub fn accept_invite(
281        &mut self,
282        invite_input: &str,
283        keys: &Keys,
284    ) -> Result<(DirectThreadSnapshot, Vec<DirectMessageCommand>), String> {
285        match self.accept_invite_with_status(invite_input, keys)? {
286            DirectInviteAcceptanceOutcome::Accepted { thread, commands } => {
287                Ok((thread, commands))
288            }
289            DirectInviteAcceptanceOutcome::PendingOwnerRoster {
290                owner_pubkey,
291                device_pubkey,
292                ..
293            } => Err(format!(
294                "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"
295            )),
296        }
297    }
298
299    pub fn accept_invite_with_status(
300        &mut self,
301        invite_input: &str,
302        _keys: &Keys,
303    ) -> Result<DirectInviteAcceptanceOutcome, String> {
304        let invite = parse_direct_invite_input(invite_input)?;
305        let owner = resolve_invite_owner(&invite, None).map_err(|error| error.to_string())?;
306        let outcome = {
307            let engine = self
308                .protocol_engine
309                .as_mut()
310                .ok_or_else(|| "Direct message runtime is not ready".to_string())?;
311            engine
312                .accept_invite(&invite, Some(owner))
313                .map_err(|error| error.to_string())?
314        };
315        let result = match outcome {
316            ProtocolAcceptInviteOutcome::Accepted(result) => result,
317            ProtocolAcceptInviteOutcome::Blocked(
318                ProtocolAcceptInviteBlock::MissingOwnerRoster {
319                    owner_pubkey,
320                    device_pubkey,
321                },
322            ) => {
323                let commands = vec![DirectMessageCommand::Subscribe {
324                    subscription_id: format!("iris-native-app-keys-{}", owner_pubkey.to_hex()),
325                    filters: vec![Filter::new()
326                        .kind(nostr::Kind::from(APP_KEYS_EVENT_KIND as u16))
327                        .author(owner_pubkey)
328                        .limit(16)],
329                    durable: true,
330                }];
331                return Ok(DirectInviteAcceptanceOutcome::PendingOwnerRoster {
332                    owner_pubkey: owner_pubkey.to_hex(),
333                    device_pubkey: device_pubkey.to_hex(),
334                    commands,
335                });
336            }
337            ProtocolAcceptInviteOutcome::Blocked(
338                ProtocolAcceptInviteBlock::UnauthorizedDevice { .. },
339            ) => {
340                return Err("Invite device is not authorized by its claimed owner".to_string());
341            }
342        };
343        let chat_id = result.owner_pubkey.to_hex();
344        self.ensure_thread(&chat_id, unix_now());
345        let mut commands = self.commands_from_effects(result.effects);
346        commands.extend(self.protocol_subscription_commands());
347        let thread = self
348            .thread(&chat_id)
349            .ok_or_else(|| "Invite chat open failed".to_string())?;
350        Ok(DirectInviteAcceptanceOutcome::Accepted { thread, commands })
351    }
352
353    pub fn send_message(
354        &mut self,
355        chat_id: &str,
356        body: &str,
357        _keys: &Keys,
358    ) -> Result<Vec<DirectMessageCommand>, String> {
359        let body = body.trim();
360        if body.is_empty() {
361            return Ok(Vec::new());
362        }
363        let public_key = PublicKey::parse(chat_id).map_err(|error| error.to_string())?;
364        let chat_id = public_key.to_hex();
365        self.ensure_thread(&chat_id, unix_now());
366        let engine = self
367            .protocol_engine
368            .as_mut()
369            .ok_or_else(|| "Direct message runtime is not ready".to_string())?;
370        let result = engine
371            .send_direct_text(public_key, &chat_id, body, None, UnixSeconds(unix_now()))
372            .map_err(|error| error.to_string())?;
373        let delivery = if result.event_ids.is_empty() {
374            DirectMessageDelivery::Pending
375        } else {
376            DirectMessageDelivery::Sent
377        };
378        self.insert_message(
379            &chat_id,
380            &result.message_id,
381            body,
382            true,
383            unix_now(),
384            delivery,
385            None,
386        );
387        Ok(self.commands_from_effects(result.effects))
388    }
389
390    pub fn process_event(&mut self, event: Event, _keys: &Keys) -> Vec<DirectMessageCommand> {
391        let event_id = event.id.to_hex();
392        if self.seen_event(&event_id) {
393            return Vec::new();
394        }
395        let Some(engine) = self.protocol_engine.as_mut() else {
396            return Vec::new();
397        };
398        let kind = event.kind.as_u16() as u32;
399        let mut effects = Vec::new();
400        let mut retry_batch = ProtocolRetryBatch::default();
401        let mut decrypted = None;
402
403        let processed = match kind {
404            APP_KEYS_EVENT_KIND if is_app_keys_event(&event) => {
405                match engine.ingest_app_keys_event(&event) {
406                    Ok(batch) => {
407                        retry_batch = batch;
408                        true
409                    }
410                    Err(error) => {
411                        self.last_error =
412                            Some(format!("Direct message device roster failed: {error}"));
413                        false
414                    }
415                }
416            }
417            INVITE_EVENT_KIND => match engine.observe_invite_event(&event) {
418                Ok(batch) => {
419                    retry_batch = batch;
420                    true
421                }
422                Err(_) => false,
423            },
424            INVITE_RESPONSE_KIND => match engine.observe_invite_response_event(&event) {
425                Ok(batch) => {
426                    retry_batch = batch;
427                    true
428                }
429                Err(_) => false,
430            },
431            MESSAGE_EVENT_KIND => match engine.process_direct_message_event(&event) {
432                Ok(message) => {
433                    decrypted = message;
434                    true
435                }
436                Err(_) => false,
437            },
438            _ => false,
439        };
440
441        if !processed {
442            return Vec::new();
443        }
444        let app_keys_owner =
445            (kind == APP_KEYS_EVENT_KIND && is_app_keys_event(&event)).then_some(event.pubkey);
446        self.mark_seen_event(&event_id);
447        if let Some(message) = decrypted {
448            self.apply_decrypted_protocol_message(message);
449        }
450        effects.extend(self.effects_from_retry_batch(retry_batch));
451        let mut commands = self.commands_from_effects(effects);
452        if app_keys_owner.is_some() {
453            commands.extend(self.protocol_subscription_commands());
454        }
455        commands
456    }
457
458    pub fn mobile_push_message_author_pubkeys(&self) -> Vec<String> {
459        let Some(engine) = self.protocol_engine.as_ref() else {
460            return Vec::new();
461        };
462        let mut authors = engine
463            .known_message_author_pubkeys()
464            .into_iter()
465            .map(|pubkey| pubkey.to_hex())
466            .collect::<Vec<_>>();
467        authors.sort();
468        authors.dedup();
469        authors
470    }
471
472    pub fn local_invite_event(&self, device_keys: &Keys) -> Option<Event> {
473        let invite = self.protocol_engine.as_ref()?.local_invite()?;
474        if invite.inviter_device_pubkey.to_bytes() != device_keys.public_key().to_bytes() {
475            return None;
476        }
477        invite_unsigned_event(&invite)
478            .ok()?
479            .sign_with_keys(device_keys)
480            .ok()
481    }
482
483    fn subscription_command(&mut self) -> Option<DirectMessageCommand> {
484        let engine = self.protocol_engine.as_ref()?;
485        let authors = engine
486            .known_message_author_pubkeys()
487            .into_iter()
488            .chain(self.owner_public_key)
489            .collect::<Vec<_>>();
490        let mut author_hexes = authors.iter().map(PublicKey::to_hex).collect::<Vec<_>>();
491        author_hexes.sort();
492        author_hexes.dedup();
493        let key = author_hexes.join(",");
494        if key.is_empty() || self.relay_subscription_key.as_deref() == Some(key.as_str()) {
495            return None;
496        }
497        self.relay_subscription_key = Some(key);
498
499        let public_keys = author_hexes
500            .iter()
501            .filter_map(|hex| PublicKey::parse(hex).ok())
502            .collect::<Vec<_>>();
503        let filter = Filter::new()
504            .authors(public_keys)
505            .kinds([
506                nostr::Kind::from(MESSAGE_EVENT_KIND as u16),
507                nostr::Kind::from(INVITE_EVENT_KIND as u16),
508                nostr::Kind::from(INVITE_RESPONSE_KIND as u16),
509                nostr::Kind::from(APP_KEYS_EVENT_KIND as u16),
510            ])
511            .limit(500);
512        Some(DirectMessageCommand::Subscribe {
513            subscription_id: "iris-native-private-chat".to_string(),
514            filters: vec![filter],
515            durable: true,
516        })
517    }
518
519    fn with_protocol_engine(self, keys: &Keys) -> Self {
520        self.with_protocol_engine_for_local_device(keys.public_key(), keys)
521    }
522
523    fn with_protocol_engine_for_local_device(
524        mut self,
525        owner: PublicKey,
526        device_keys: &Keys,
527    ) -> Self {
528        let owner_hex = owner.to_hex();
529        let device_hex = device_keys.public_key().to_hex();
530        let storage = Arc::new(SqliteStorageAdapter::new(
531            Arc::clone(&self.conn),
532            owner_hex.clone(),
533            device_hex,
534        ));
535        match ProtocolEngine::load_or_create_for_local_device(storage, owner, device_keys) {
536            Ok(engine) => {
537                self.protocol_engine = Some(engine);
538                self.owner_public_key = Some(owner);
539            }
540            Err(error) => self.last_error = Some(format!("Direct message init failed: {error}")),
541        }
542        self
543    }
544
545    fn protocol_subscription_commands(&mut self) -> Vec<DirectMessageCommand> {
546        self.subscription_command().into_iter().collect()
547    }
548
549    fn commands_from_effects(&mut self, effects: Vec<ProtocolEffect>) -> Vec<DirectMessageCommand> {
550        let mut commands = Vec::new();
551        for effect in effects {
552            match effect {
553                ProtocolEffect::Publish(publish) => {
554                    commands.push(DirectMessageCommand::Publish(publish.event));
555                }
556            }
557        }
558        commands
559    }
560
561    fn effects_from_retry_batch(&mut self, batch: ProtocolRetryBatch) -> Vec<ProtocolEffect> {
562        let mut effects = batch.effects;
563        effects.extend(batch.group_result.effects);
564        for message in batch.direct_messages {
565            self.apply_decrypted_protocol_message(message);
566        }
567        effects
568    }
569
570    fn apply_decrypted_protocol_message(&mut self, message: ProtocolDecryptedMessage) {
571        self.apply_decrypted(
572            message.sender,
573            message.conversation_owner,
574            &message.content,
575            message.event_id,
576        );
577    }
578
579    fn apply_decrypted(
580        &mut self,
581        sender: PublicKey,
582        conversation_owner: Option<PublicKey>,
583        content: &str,
584        source_event_id: Option<String>,
585    ) {
586        let Some(rumor) = parse_runtime_rumor(content) else {
587            return;
588        };
589        if rumor.kind != CHAT_MESSAGE_KIND {
590            return;
591        }
592        let local_owner = self.owner_public_key;
593        let peer = if local_owner == Some(sender) {
594            conversation_owner.unwrap_or(sender)
595        } else {
596            sender
597        };
598        let chat_id = peer.to_hex();
599        self.ensure_thread(&chat_id, rumor.created_at_secs);
600        self.insert_message(
601            &chat_id,
602            &rumor.id,
603            &rumor.content,
604            local_owner == Some(sender),
605            rumor.created_at_secs,
606            if local_owner == Some(sender) {
607                DirectMessageDelivery::Sent
608            } else {
609                DirectMessageDelivery::Received
610            },
611            source_event_id.as_deref(),
612        );
613    }
614
615    fn ensure_schema(&self) {
616        if let Ok(conn) = self.conn.lock() {
617            let _ = conn.execute_batch(SCHEMA);
618        }
619    }
620
621    fn ensure_thread(&self, chat_id: &str, updated_at: u64) {
622        if let Ok(conn) = self.conn.lock() {
623            let _ = conn.execute(
624                "INSERT INTO private_chat_threads (chat_id, display_name, avatar_seed, updated_at_secs)
625                 VALUES (?1, '', '', ?2)
626                 ON CONFLICT(chat_id) DO UPDATE SET updated_at_secs = MAX(updated_at_secs, excluded.updated_at_secs)",
627                params![chat_id, updated_at as i64],
628            );
629        }
630    }
631
632    fn messages(&self, chat_id: &str, limit: usize) -> Vec<DirectMessageSnapshot> {
633        let Ok(conn) = self.conn.lock() else {
634            return Vec::new();
635        };
636        let mut stmt = match conn.prepare(
637            "SELECT id, body, is_outgoing, created_at_secs, delivery
638             FROM private_chat_messages
639             WHERE chat_id = ?1
640             ORDER BY created_at_secs DESC, id DESC
641             LIMIT ?2",
642        ) {
643            Ok(stmt) => stmt,
644            Err(_) => return Vec::new(),
645        };
646        let rows = match stmt.query_map(params![chat_id, limit as i64], |row| {
647            Ok(DirectMessageSnapshot {
648                id: row.get(0)?,
649                chat_id: chat_id.to_string(),
650                body: row.get(1)?,
651                is_outgoing: row.get::<_, i64>(2)? != 0,
652                created_at_secs: row.get::<_, i64>(3)?.max(0) as u64,
653                delivery: DirectMessageDelivery::from_str(&row.get::<_, String>(4)?),
654            })
655        }) {
656            Ok(rows) => rows,
657            Err(_) => return Vec::new(),
658        };
659        let mut messages = rows.filter_map(Result::ok).collect::<Vec<_>>();
660        messages.reverse();
661        messages
662    }
663
664    #[allow(clippy::too_many_arguments)]
665    fn insert_message(
666        &self,
667        chat_id: &str,
668        id: &str,
669        body: &str,
670        is_outgoing: bool,
671        created_at: u64,
672        delivery: DirectMessageDelivery,
673        source_event_id: Option<&str>,
674    ) {
675        if id.is_empty() {
676            return;
677        }
678        if let Ok(conn) = self.conn.lock() {
679            let _ = conn.execute(
680                "INSERT OR IGNORE INTO private_chat_messages
681                 (chat_id, id, body, is_outgoing, created_at_secs, delivery, source_event_id)
682                 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
683                params![
684                    chat_id,
685                    id,
686                    body,
687                    is_outgoing as i64,
688                    created_at as i64,
689                    delivery.as_str(),
690                    source_event_id,
691                ],
692            );
693            let _ = conn.execute(
694                "UPDATE private_chat_threads SET updated_at_secs = MAX(updated_at_secs, ?2)
695                 WHERE chat_id = ?1",
696                params![chat_id, created_at as i64],
697            );
698        }
699    }
700
701    fn seen_event(&self, event_id: &str) -> bool {
702        let Ok(conn) = self.conn.lock() else {
703            return true;
704        };
705        conn.query_row(
706            "SELECT 1 FROM private_chat_seen_events WHERE event_id = ?1",
707            [event_id],
708            |_| Ok(()),
709        )
710        .optional()
711        .ok()
712        .flatten()
713        .is_some()
714    }
715
716    fn mark_seen_event(&self, event_id: &str) {
717        if let Ok(conn) = self.conn.lock() {
718            let _ = conn.execute(
719                "INSERT OR IGNORE INTO private_chat_seen_events (event_id) VALUES (?1)",
720                [event_id],
721            );
722        }
723    }
724}
725
726struct RuntimeRumor {
727    id: String,
728    kind: u32,
729    content: String,
730    created_at_secs: u64,
731}
732
733fn parse_runtime_rumor(content: &str) -> Option<RuntimeRumor> {
734    let mut event = serde_json::from_str::<UnsignedEvent>(content).ok()?;
735    event.ensure_id();
736    event.verify_id().ok()?;
737    Some(RuntimeRumor {
738        id: event.id.as_ref()?.to_string(),
739        kind: event.kind.as_u16() as u32,
740        content: event.content,
741        created_at_secs: event.created_at.as_secs(),
742    })
743}
744
745fn chat_snapshot_for_pubkey(chat_id: &str) -> DirectChatSnapshot {
746    DirectChatSnapshot {
747        chat_id: chat_id.to_string(),
748        last_message_preview: String::new(),
749        last_message_at: 0,
750        unread_count: 0,
751    }
752}
753
754fn normalize_pubkey(input: &str) -> Result<String, String> {
755    PublicKey::parse(input)
756        .map(|pubkey| pubkey.to_hex())
757        .map_err(|error| error.to_string())
758}
759
760fn parse_direct_invite_input(input: &str) -> Result<Invite, String> {
761    let trimmed = input.trim();
762    if trimmed.is_empty() {
763        return Err("Invite link is required".to_string());
764    }
765    if let Ok(invite) = parse_invite_url(trimmed) {
766        return Ok(invite);
767    }
768
769    let mut candidates = vec![trimmed.to_string()];
770    if let Some((_, fragment)) = trimmed.split_once('#') {
771        candidates.push(fragment.to_string());
772        candidates.push(fragment.trim_start_matches('/').to_string());
773        candidates.extend(
774            fragment
775                .split(['/', '?', '&', '='])
776                .filter(|part| !part.trim().is_empty())
777                .map(ToString::to_string),
778        );
779    }
780    if let Some((_, query)) = trimmed.split_once('?') {
781        candidates.extend(
782            query
783                .split(['/', '?', '&', '='])
784                .filter(|part| !part.trim().is_empty())
785                .map(ToString::to_string),
786        );
787    }
788
789    for candidate in candidates {
790        let candidate = candidate.trim().trim_start_matches('/');
791        let candidate = candidate.strip_prefix("invite/").unwrap_or(candidate);
792        if candidate.is_empty() || candidate.eq_ignore_ascii_case("invite") {
793            continue;
794        }
795        for wrapped in [
796            candidate.to_string(),
797            format!("https://chat.iris.to#{candidate}"),
798            format!("https://chat.iris.to#/{candidate}"),
799        ] {
800            if let Ok(invite) = parse_invite_url(&wrapped) {
801                return Ok(invite);
802            }
803        }
804    }
805
806    parse_invite_url(trimmed).map_err(|error| error.to_string())
807}
808
809fn unix_now() -> u64 {
810    std::time::SystemTime::now()
811        .duration_since(std::time::UNIX_EPOCH)
812        .map(|duration| duration.as_secs())
813        .unwrap_or_default()
814}
815
816#[cfg(test)]
817mod tests {
818    use super::*;
819    use crate::{invite_url, parse_invite_event, DeviceEntry};
820    use nostr::Kind;
821
822    fn publish_events(commands: Vec<DirectMessageCommand>) -> Vec<Event> {
823        commands
824            .into_iter()
825            .filter_map(|command| match command {
826                DirectMessageCommand::Publish(event) => Some(event),
827                DirectMessageCommand::Subscribe { .. } => None,
828            })
829            .collect()
830    }
831
832    fn publish_kinds(commands: &[DirectMessageCommand]) -> Vec<Kind> {
833        commands
834            .iter()
835            .filter_map(|command| match command {
836                DirectMessageCommand::Publish(event) => Some(event.kind),
837                DirectMessageCommand::Subscribe { .. } => None,
838            })
839            .collect()
840    }
841
842    fn route_wrapped_invite_url(invite: &Invite) -> String {
843        let raw = invite_url(invite, "https://chat.iris.to").expect("invite url");
844        let Some((_, fragment)) = raw.split_once('#') else {
845            return raw;
846        };
847        let payload = fragment.trim_start_matches('/');
848        if payload.starts_with("invite/") {
849            raw
850        } else {
851            format!("https://chat.iris.to/#/invite/{payload}")
852        }
853    }
854
855    #[test]
856    fn accepts_route_wrapped_invite_and_sends_direct_message() {
857        let inviter_keys = Keys::generate();
858        let accepter_keys = Keys::generate();
859        let mut inviter =
860            DirectMessageService::memory_for_local_device(inviter_keys.public_key(), &inviter_keys);
861        let mut accepter = DirectMessageService::memory_for_local_device(
862            accepter_keys.public_key(),
863            &accepter_keys,
864        );
865        let invite_event = inviter
866            .local_invite_event(&inviter_keys)
867            .expect("local invite event");
868        let invite = parse_invite_event(&invite_event).expect("invite event");
869        let invite_url = route_wrapped_invite_url(&invite);
870
871        let (thread, accept_commands) = accepter
872            .accept_invite(&invite_url, &accepter_keys)
873            .expect("accept invite");
874        assert_eq!(thread.chat.chat_id, inviter_keys.public_key().to_hex());
875        let accept_kinds = publish_kinds(&accept_commands);
876        assert!(accept_kinds.contains(&Kind::from(INVITE_RESPONSE_KIND as u16)));
877        assert!(accept_kinds.contains(&Kind::from(MESSAGE_EVENT_KIND as u16)));
878
879        for event in publish_events(accept_commands) {
880            inviter.process_event(event, &inviter_keys);
881        }
882
883        let send_commands = accepter
884            .send_message(
885                &inviter_keys.public_key().to_hex(),
886                "hello from invite accepter",
887                &accepter_keys,
888            )
889            .expect("send message");
890        assert!(publish_kinds(&send_commands).contains(&Kind::from(MESSAGE_EVENT_KIND as u16)));
891
892        for event in publish_events(send_commands) {
893            inviter.process_event(event, &inviter_keys);
894        }
895
896        let inviter_thread = inviter
897            .thread(&accepter_keys.public_key().to_hex())
898            .expect("inviter thread");
899        assert_eq!(inviter_thread.messages.len(), 1);
900        assert_eq!(
901            inviter_thread.messages[0].body,
902            "hello from invite accepter"
903        );
904        assert!(!inviter_thread.messages[0].is_outgoing);
905        assert_eq!(
906            inviter_thread.messages[0].delivery,
907            DirectMessageDelivery::Received
908        );
909    }
910
911    #[test]
912    fn claimed_owner_invite_retries_after_ordinary_app_keys_ingestion() {
913        let inviter_owner = Keys::generate();
914        let inviter_device = Keys::generate();
915        let accepter_keys = Keys::generate();
916        let mut accepter = DirectMessageService::memory_for_local_device(
917            accepter_keys.public_key(),
918            &accepter_keys,
919        );
920        let mut invite = Invite::create_new(
921            inviter_device.public_key(),
922            Some(inviter_device.public_key().to_hex()),
923            Some(1),
924        )
925        .expect("invite");
926        invite.owner_public_key = Some(inviter_owner.public_key());
927        invite.purpose = Some("private".to_string());
928        let invite_url = route_wrapped_invite_url(&invite);
929
930        let pending = accepter
931            .accept_invite_with_status(&invite_url, &accepter_keys)
932            .expect("pending acceptance");
933        let commands = match pending {
934            DirectInviteAcceptanceOutcome::PendingOwnerRoster {
935                owner_pubkey,
936                device_pubkey,
937                commands,
938            } => {
939                assert_eq!(owner_pubkey, inviter_owner.public_key().to_hex());
940                assert_eq!(device_pubkey, inviter_device.public_key().to_hex());
941                commands
942            }
943            DirectInviteAcceptanceOutcome::Accepted { .. } => {
944                panic!("owner claim must wait for AppKeys")
945            }
946        };
947        let owner_hex = inviter_owner.public_key().to_hex();
948        assert!(commands.iter().any(|command| {
949            matches!(
950                command, DirectMessageCommand::Subscribe { filters, durable: true, .. }
951                    if filters.iter().any(|filter| {
952                        serde_json::to_string(filter)
953                            .is_ok_and(|json| json.contains(&owner_hex))
954                    })
955            )
956        }));
957        assert!(accepter
958            .chats()
959            .iter()
960            .all(|chat| chat.chat_id != owner_hex));
961
962        let roster = AppKeys::new(vec![DeviceEntry::new(inviter_device.public_key(), 10)])
963            .get_event_at(inviter_owner.public_key(), 10)
964            .sign_with_keys(&inviter_owner)
965            .expect("signed inviter AppKeys");
966        let completion = accepter.process_event(roster.clone(), &accepter_keys);
967
968        assert!(accepter
969            .chats()
970            .iter()
971            .all(|chat| chat.chat_id != owner_hex));
972        assert!(publish_kinds(&completion).is_empty());
973        let accepted = accepter
974            .accept_invite_with_status(&invite_url, &accepter_keys)
975            .expect("authorized retry");
976        assert!(matches!(
977            accepted,
978            DirectInviteAcceptanceOutcome::Accepted { .. }
979        ));
980    }
981}