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}