1use super::ConnectionExt;
2use super::schema::conversation_list::dsl::conversation_list;
3use crate::consent_record::ConsentState;
4use crate::group::{ConversationType, GroupMembershipState, GroupQueryArgs, GroupQueryOrderBy};
5use crate::group_message::{ContentType, DeliveryStatus, GroupMessageKind};
6use crate::{DbConnection, StorageError};
7use diesel::dsl::sql;
8use diesel::{
9 BoolExpressionMethods, ExpressionMethods, JoinOnDsl, QueryDsl, Queryable, RunQueryDsl, Table,
10};
11use serde::{Deserialize, Serialize};
12
13#[derive(Queryable, Debug, Clone, Deserialize, Serialize)]
14#[diesel(table_name = conversation_list)]
15#[diesel(primary_key(id))]
16pub struct ConversationListItem {
18 pub id: xmtp_proto::types::GroupId,
20 pub created_at_ns: i64,
22 pub membership_state: GroupMembershipState,
24 pub installations_last_checked: i64,
26 pub added_by_inbox_id: String,
28 pub welcome_sequence_id: Option<i64>,
30 pub dm_id: Option<String>,
32 pub rotated_at_ns: i64,
34 pub conversation_type: ConversationType,
36 pub is_commit_log_forked: Option<bool>,
38 pub message_id: Option<Vec<u8>>,
40 pub decrypted_message_bytes: Option<Vec<u8>>,
42 pub sent_at_ns: Option<i64>,
44 pub kind: Option<GroupMessageKind>,
46 pub sender_installation_id: Option<Vec<u8>>,
48 pub sender_inbox_id: Option<String>,
50 pub delivery_status: Option<DeliveryStatus>,
52 pub content_type: Option<ContentType>,
54 pub version_major: Option<i32>,
56 pub version_minor: Option<i32>,
58 pub authority_id: Option<String>,
60 pub sequence_id: Option<i64>,
62}
63
64pub trait QueryConversationList {
65 fn fetch_conversation_list<A: AsRef<GroupQueryArgs>>(
66 &self,
67 args: A,
68 ) -> Result<Vec<ConversationListItem>, StorageError>;
69}
70
71impl<T> QueryConversationList for &T
72where
73 T: QueryConversationList,
74{
75 fn fetch_conversation_list<A: AsRef<GroupQueryArgs>>(
76 &self,
77 args: A,
78 ) -> Result<Vec<ConversationListItem>, StorageError> {
79 (**self).fetch_conversation_list(args)
80 }
81}
82
83impl<C: ConnectionExt> QueryConversationList for DbConnection<C> {
84 fn fetch_conversation_list<A: AsRef<GroupQueryArgs>>(
85 &self,
86 args: A,
87 ) -> Result<Vec<ConversationListItem>, StorageError> {
88 use crate::schema::consent_records::dsl as consent_dsl;
89 use crate::schema::conversation_list::dsl as conversation_list_dsl;
90
91 args.as_ref().validate()?;
92
93 let GroupQueryArgs {
94 allowed_states,
95 created_after_ns,
96 created_before_ns,
97 limit,
98 conversation_type,
99 consent_states,
100 include_sync_groups,
101 include_duplicate_dms,
102 last_activity_after_ns,
103 last_activity_before_ns,
104 order_by,
105 ..
106 } = args.as_ref();
107
108 let order_expression = match order_by.clone().unwrap_or_default() {
109 GroupQueryOrderBy::CreatedAt => {
110 diesel::dsl::sql::<diesel::sql_types::BigInt>("created_at_ns DESC")
111 }
112 GroupQueryOrderBy::LastActivity => diesel::dsl::sql::<diesel::sql_types::BigInt>(
113 "COALESCE(sent_at_ns, created_at_ns) DESC",
114 ),
115 };
116
117 let mut query = conversation_list
118 .select(conversation_list::all_columns())
119 .filter(
120 conversation_list_dsl::conversation_type.ne_all(ConversationType::virtual_types()),
121 )
122 .order(order_expression)
123 .into_boxed();
124
125 if !include_duplicate_dms {
126 query = query.filter(sql::<diesel::sql_types::Bool>(
129 "NOT EXISTS (
130 SELECT 1 FROM groups g2
131 WHERE COALESCE(g2.dm_id, g2.id) = COALESCE(conversation_list.dm_id, conversation_list.id)
132 AND (COALESCE(g2.last_message_ns, 0) > COALESCE((
133 SELECT g1.last_message_ns FROM groups g1 WHERE g1.id = conversation_list.id
134 ), 0)
135 OR (COALESCE(g2.last_message_ns, 0) = COALESCE((
136 SELECT g1.last_message_ns FROM groups g1 WHERE g1.id = conversation_list.id
137 ), 0) AND g2.id > conversation_list.id))
138 )",
139 ));
140 }
141
142 if let Some(limit) = limit {
143 query = query.limit(*limit);
144 }
145
146 if let Some(allowed_states) = allowed_states {
147 query = query.filter(conversation_list_dsl::membership_state.eq_any(allowed_states));
148 }
149
150 if let Some(last_activity_after_ns) = last_activity_after_ns {
152 query = query.filter(
155 diesel::dsl::sql::<diesel::sql_types::BigInt>(
156 "COALESCE(sent_at_ns, created_at_ns)",
157 )
158 .gt(last_activity_after_ns),
159 );
160 }
161
162 if let Some(created_after_ns) = created_after_ns {
163 query = query.filter(conversation_list_dsl::created_at_ns.gt(created_after_ns));
164 }
165
166 if let Some(last_activity_before_ns) = last_activity_before_ns {
167 query = query.filter(
168 diesel::dsl::sql::<diesel::sql_types::BigInt>(
169 "COALESCE(sent_at_ns, created_at_ns)",
170 )
171 .lt(last_activity_before_ns),
172 );
173 }
174
175 if let Some(created_before_ns) = created_before_ns {
176 query = query.filter(conversation_list_dsl::created_at_ns.lt(created_before_ns));
177 }
178
179 if let Some(conversation_type) = conversation_type {
180 query = query.filter(conversation_list_dsl::conversation_type.eq(conversation_type));
181 }
182
183 let effective_consent_states = match &consent_states {
184 Some(states) if !states.is_empty() => states.clone(),
185 _ => vec![ConsentState::Allowed, ConsentState::Unknown],
186 };
187
188 let includes_unknown = effective_consent_states.contains(&ConsentState::Unknown);
189 let includes_all = effective_consent_states.len() == 3;
190
191 let filtered_states: Vec<_> = effective_consent_states
192 .iter()
193 .filter(|state| **state != ConsentState::Unknown)
194 .cloned()
195 .collect();
196
197 let mut conversations = if includes_all {
198 self.raw_query(|conn| query.load::<ConversationListItem>(conn))?
200 } else if includes_unknown {
201 let left_joined_query = query
203 .left_join(
204 consent_dsl::consent_records.on(sql::<diesel::sql_types::Text>(
205 "lower(hex(conversation_list.id))",
206 )
207 .eq(consent_dsl::entity)),
208 )
209 .filter(
210 consent_dsl::state
211 .is_null()
212 .or(consent_dsl::state.eq(ConsentState::Unknown))
213 .or(consent_dsl::state.eq_any(filtered_states.clone())),
214 )
215 .select(conversation_list::all_columns());
216
217 self.raw_query(|conn| left_joined_query.load::<ConversationListItem>(conn))?
218 } else {
219 let inner_joined_query = query
221 .inner_join(
222 consent_dsl::consent_records.on(sql::<diesel::sql_types::Text>(
223 "lower(hex(conversation_list.id))",
224 )
225 .eq(consent_dsl::entity)),
226 )
227 .filter(consent_dsl::state.eq_any(filtered_states.clone()))
228 .select(conversation_list::all_columns());
229
230 self.raw_query(|conn| inner_joined_query.load::<ConversationListItem>(conn))?
231 };
232
233 if matches!(conversation_type, Some(ConversationType::Sync)) || *include_sync_groups {
236 let query = conversation_list_dsl::conversation_list
237 .filter(conversation_list_dsl::conversation_type.eq(ConversationType::Sync));
238 let mut sync_groups = self.raw_query(|conn| query.load(conn))?;
239 conversations.append(&mut sync_groups);
240 }
241
242 Ok(conversations)
243 }
244}
245
246#[cfg(test)]
247pub(crate) mod tests {
248 use crate::Store;
249 use crate::consent_record::{ConsentState, ConsentType};
250 use crate::group::tests::{
251 generate_consent_record, generate_dm, generate_group, generate_group_with_created_at,
252 };
253 use crate::group::{GroupMembershipState, GroupQueryArgs, GroupQueryOrderBy};
254 use crate::group_message::ContentType;
255 use crate::group_message::tests::generate_message;
256 use crate::prelude::*;
257 use crate::test_utils::with_connection;
258
259 #[xmtp_common::test]
260 fn test_single_group_multiple_messages() {
261 with_connection(|conn| {
262 let group = generate_group(None);
264 group.store(conn).unwrap();
265
266 for i in 1..5 {
268 let message = crate::encrypted_store::group_message::tests::generate_message(
269 None,
270 Some(&group.id),
271 Some(i * 1000),
272 Some(ContentType::Text),
273 None,
274 None,
275 );
276
277 message.store(conn).unwrap();
278 }
279
280 let conversation_list = conn
282 .fetch_conversation_list(GroupQueryArgs::default())
283 .unwrap();
284 assert_eq!(conversation_list.len(), 1, "Should return one group");
285 assert_eq!(
286 conversation_list[0].id, group.id,
287 "Returned group ID should match the created group"
288 );
289 assert_eq!(
290 conversation_list[0].sent_at_ns.unwrap(),
291 4000,
292 "Last message should be the most recent one"
293 );
294 })
295 }
296
297 #[xmtp_common::test]
298 fn test_three_groups_specific_ordering() {
299 with_connection(|conn| {
300 let group_a = generate_group_with_created_at(None, 5000); let group_b = generate_group_with_created_at(None, 2000); let group_c = generate_group_with_created_at(None, 1000); group_a.store(conn).unwrap();
306 group_b.store(conn).unwrap();
307 group_c.store(conn).unwrap();
308 let message = crate::encrypted_store::group_message::tests::generate_message(
310 None,
311 Some(&group_b.id),
312 Some(3000), None,
314 None,
315 None,
316 );
317 message.store(conn).unwrap();
318
319 let conversation_list = conn
321 .fetch_conversation_list(GroupQueryArgs::default())
322 .unwrap();
323
324 assert_eq!(conversation_list.len(), 3, "Should return all three groups");
325 assert_eq!(
326 conversation_list[0].id, group_a.id,
327 "Group created after the last message should come first"
328 );
329 assert_eq!(
330 conversation_list[1].id, group_b.id,
331 "Group with the last message should come second"
332 );
333 assert_eq!(
334 conversation_list[2].id, group_c.id,
335 "Group created before the last message with no messages should come last"
336 );
337 })
338 }
339
340 #[xmtp_common::test]
341 fn test_group_with_newer_message_update() {
342 with_connection(|conn| {
343 let group = generate_group(None);
345 group.store(conn).unwrap();
346
347 let first_message = crate::encrypted_store::group_message::tests::generate_message(
349 None,
350 Some(&group.id),
351 Some(1000),
352 Some(ContentType::Text),
353 None,
354 None,
355 );
356 first_message.store(conn).unwrap();
357
358 let mut conversation_list = conn
360 .fetch_conversation_list(GroupQueryArgs::default())
361 .unwrap();
362 assert_eq!(conversation_list.len(), 1, "Should return one group");
363 assert_eq!(
364 conversation_list[0].sent_at_ns.unwrap(),
365 1000,
366 "Last message should match the first message"
367 );
368
369 let second_message = crate::encrypted_store::group_message::tests::generate_message(
371 None,
372 Some(&group.id),
373 Some(2000),
374 Some(ContentType::Text),
375 None,
376 None,
377 );
378 second_message.store(conn).unwrap();
379
380 conversation_list = conn
382 .fetch_conversation_list(GroupQueryArgs::default())
383 .unwrap();
384 assert_eq!(
385 conversation_list[0].sent_at_ns.unwrap(),
386 2000,
387 "Last message should now match the second (newest) message"
388 );
389 })
390 }
391
392 #[xmtp_common::test]
393 fn test_find_conversations_by_consent_state() {
394 with_connection(|conn| {
395 let test_group_1 = generate_group(Some(GroupMembershipState::Allowed));
396 test_group_1.store(conn).unwrap();
397 let test_group_2 = generate_group(Some(GroupMembershipState::Allowed));
398 test_group_2.store(conn).unwrap();
399 let test_group_3 = generate_dm(Some(GroupMembershipState::Allowed));
400 test_group_3.store(conn).unwrap();
401 let test_group_4 = generate_dm(Some(GroupMembershipState::Allowed));
402 test_group_4.store(conn).unwrap();
403
404 let test_group_1_consent = generate_consent_record(
405 ConsentType::ConversationId,
406 ConsentState::Allowed,
407 hex::encode(test_group_1.id),
408 );
409 test_group_1_consent.store(conn).unwrap();
410 let test_group_2_consent = generate_consent_record(
411 ConsentType::ConversationId,
412 ConsentState::Denied,
413 hex::encode(test_group_2.id),
414 );
415 test_group_2_consent.store(conn).unwrap();
416 let test_group_3_consent = generate_consent_record(
417 ConsentType::ConversationId,
418 ConsentState::Allowed,
419 hex::encode(test_group_3.id),
420 );
421 test_group_3_consent.store(conn).unwrap();
422
423 let all_results = conn
424 .fetch_conversation_list(GroupQueryArgs {
425 consent_states: Some(vec![
426 ConsentState::Allowed,
427 ConsentState::Unknown,
428 ConsentState::Denied,
429 ]),
430 ..Default::default()
431 })
432 .unwrap();
433 assert_eq!(all_results.len(), 4);
434
435 let default_results = conn
436 .fetch_conversation_list(GroupQueryArgs::default())
437 .unwrap();
438 assert_eq!(default_results.len(), 3);
439
440 let allowed_results = conn
441 .fetch_conversation_list(GroupQueryArgs {
442 consent_states: Some(vec![ConsentState::Allowed]),
443 ..Default::default()
444 })
445 .unwrap();
446 assert_eq!(allowed_results.len(), 2);
447
448 let allowed_unknown_results = conn
449 .fetch_conversation_list(GroupQueryArgs {
450 consent_states: Some(vec![ConsentState::Allowed, ConsentState::Unknown]),
451 ..Default::default()
452 })
453 .unwrap();
454 assert_eq!(allowed_unknown_results.len(), 3);
455
456 let denied_results = conn
457 .fetch_conversation_list(GroupQueryArgs {
458 consent_states: Some(vec![ConsentState::Denied]),
459 ..Default::default()
460 })
461 .unwrap();
462 assert_eq!(denied_results.len(), 1);
463 assert_eq!(denied_results[0].id, test_group_2.id);
464
465 let unknown_results = conn
466 .fetch_conversation_list(GroupQueryArgs {
467 consent_states: Some(vec![ConsentState::Unknown]),
468 ..Default::default()
469 })
470 .unwrap();
471 assert_eq!(unknown_results.len(), 1);
472 assert_eq!(unknown_results[0].id, test_group_4.id);
473
474 let empty_array_results = conn
475 .fetch_conversation_list(GroupQueryArgs {
476 consent_states: Some(vec![]),
477 ..Default::default()
478 })
479 .unwrap();
480 assert_eq!(empty_array_results.len(), 3);
481 })
482 }
483
484 #[xmtp_common::test]
485 fn test_find_conversations_default_excludes_denied() {
486 with_connection(|conn| {
487 let allowed_group = generate_group(Some(GroupMembershipState::Allowed));
489 allowed_group.store(conn).unwrap();
490
491 let denied_group = generate_group(Some(GroupMembershipState::Allowed));
492 denied_group.store(conn).unwrap();
493
494 let unknown_group = generate_group(Some(GroupMembershipState::Allowed));
495 unknown_group.store(conn).unwrap();
496
497 let allowed_consent = generate_consent_record(
499 ConsentType::ConversationId,
500 ConsentState::Allowed,
501 hex::encode(allowed_group.id),
502 );
503 allowed_consent.store(conn).unwrap();
504
505 let denied_consent = generate_consent_record(
506 ConsentType::ConversationId,
507 ConsentState::Denied,
508 hex::encode(denied_group.id),
509 );
510 denied_consent.store(conn).unwrap();
511
512 let default_results = conn
514 .fetch_conversation_list(GroupQueryArgs::default())
515 .unwrap();
516
517 assert_eq!(default_results.len(), 2);
519 let returned_ids: Vec<_> = default_results.iter().map(|g| &g.id).collect();
520 assert!(returned_ids.contains(&&allowed_group.id));
521 assert!(returned_ids.contains(&&unknown_group.id));
522 assert!(!returned_ids.contains(&&denied_group.id));
523 })
524 }
525
526 #[xmtp_common::test(unwrap_try = true)]
527 fn test_unknown_content_type_is_present() {
528 with_connection(|conn| {
529 let dm = generate_dm(None);
530 dm.store(conn)?;
531
532 let m = generate_message(
533 None,
534 Some(&dm.id),
535 Some(5000),
536 Some(ContentType::Unknown),
537 None,
538 None,
539 );
540 m.store(conn)?;
541
542 let conv = conn.fetch_conversation_list(GroupQueryArgs {
543 ..Default::default()
544 })?;
545
546 assert!(conv[0].message_id.is_some());
548 })
549 }
550
551 #[xmtp_common::test]
552 fn test_last_activity_after_ns_filter() {
553 with_connection(|conn| {
554 let group1 = generate_group_with_created_at(None, 1000);
556 let group2 = generate_group_with_created_at(None, 2000);
557 let group3 = generate_group_with_created_at(None, 3000);
558
559 group1.store(conn).unwrap();
560 group2.store(conn).unwrap();
561 group3.store(conn).unwrap();
562
563 let message1 = crate::encrypted_store::group_message::tests::generate_message(
565 None,
566 Some(&group1.id),
567 Some(5000),
568 Some(ContentType::Text),
569 None,
570 None,
571 );
572 message1.store(conn).unwrap();
573
574 let message2 = crate::encrypted_store::group_message::tests::generate_message(
576 None,
577 Some(&group2.id),
578 Some(4000),
579 Some(ContentType::Text),
580 None,
581 None,
582 );
583 message2.store(conn).unwrap();
584
585 let results = conn
589 .fetch_conversation_list(GroupQueryArgs {
590 last_activity_after_ns: Some(3500),
591 ..Default::default()
592 })
593 .unwrap();
594 assert_eq!(
595 results.len(),
596 2,
597 "Should return groups with activity after 3500"
598 );
599
600 let returned_ids: Vec<_> = results.iter().map(|g| &g.id).collect();
601 assert!(
602 returned_ids.contains(&&group1.id),
603 "Should include group1 (message at 5000)"
604 );
605 assert!(
606 returned_ids.contains(&&group2.id),
607 "Should include group2 (message at 4000)"
608 );
609 assert!(
610 !returned_ids.contains(&&group3.id),
611 "Should not include group3 (created at 3000)"
612 );
613
614 let results = conn
616 .fetch_conversation_list(GroupQueryArgs {
617 last_activity_after_ns: Some(4500),
618 ..Default::default()
619 })
620 .unwrap();
621 assert_eq!(results.len(), 1, "Should return only group1");
622 assert_eq!(results[0].id, group1.id, "Should be group1");
623
624 let results = conn
626 .fetch_conversation_list(GroupQueryArgs {
627 last_activity_after_ns: Some(2500),
628 ..Default::default()
629 })
630 .unwrap();
631 assert_eq!(results.len(), 3, "Should return all groups");
632 })
633 }
634
635 #[xmtp_common::test]
636 fn test_last_activity_before_ns_filter() {
637 with_connection(|conn| {
638 let group1 = generate_group_with_created_at(None, 1000);
640 let group2 = generate_group_with_created_at(None, 2000);
641 let group3 = generate_group_with_created_at(None, 3000);
642
643 group1.store(conn).unwrap();
644 group2.store(conn).unwrap();
645 group3.store(conn).unwrap();
646
647 let message1 = crate::encrypted_store::group_message::tests::generate_message(
649 None,
650 Some(&group1.id),
651 Some(5000),
652 Some(ContentType::Text),
653 None,
654 None,
655 );
656 message1.store(conn).unwrap();
657
658 let message2 = crate::encrypted_store::group_message::tests::generate_message(
660 None,
661 Some(&group2.id),
662 Some(4000),
663 Some(ContentType::Text),
664 None,
665 None,
666 );
667 message2.store(conn).unwrap();
668
669 let results = conn
673 .fetch_conversation_list(GroupQueryArgs {
674 last_activity_before_ns: Some(4500),
675 ..Default::default()
676 })
677 .unwrap();
678 assert_eq!(
679 results.len(),
680 2,
681 "Should return groups with activity before 4500"
682 );
683
684 let returned_ids: Vec<_> = results.iter().map(|g| &g.id).collect();
685 assert!(
686 !returned_ids.contains(&&group1.id),
687 "Should not include group1 (message at 5000)"
688 );
689 assert!(
690 returned_ids.contains(&&group2.id),
691 "Should include group2 (message at 4000)"
692 );
693 assert!(
694 returned_ids.contains(&&group3.id),
695 "Should include group3 (created at 3000)"
696 );
697
698 let results = conn
700 .fetch_conversation_list(GroupQueryArgs {
701 last_activity_before_ns: Some(3500),
702 ..Default::default()
703 })
704 .unwrap();
705 assert_eq!(results.len(), 1, "Should return only group3");
706 assert_eq!(results[0].id, group3.id, "Should be group3");
707
708 let results = conn
710 .fetch_conversation_list(GroupQueryArgs {
711 last_activity_before_ns: Some(5500),
712 ..Default::default()
713 })
714 .unwrap();
715 assert_eq!(results.len(), 3, "Should return all groups");
716 })
717 }
718
719 #[xmtp_common::test]
720 fn test_activity_filters_combined_with_limit() {
721 with_connection(|conn| {
722 let mut groups = Vec::new();
724 for i in 0..5 {
725 let group = generate_group_with_created_at(None, (i + 1) * 1000);
726 group.store(conn).unwrap();
727
728 let message = crate::encrypted_store::group_message::tests::generate_message(
730 None,
731 Some(&group.id),
732 Some((100 - i) * 1000), Some(ContentType::Text),
734 None,
735 None,
736 );
737 message.store(conn).unwrap();
738 groups.push(group);
739 }
740
741 let results = conn
744 .fetch_conversation_list(GroupQueryArgs {
745 last_activity_after_ns: Some(96_000),
746 limit: Some(2),
747 order_by: Some(GroupQueryOrderBy::LastActivity),
748 ..Default::default()
749 })
750 .unwrap();
751 assert_eq!(results.len(), 2, "Should return 2 groups due to limit");
752
753 assert_eq!(
755 results[0].sent_at_ns.unwrap(),
756 100_000,
757 "First should be most recent"
758 );
759 assert_eq!(
760 results[1].sent_at_ns.unwrap(),
761 99_000,
762 "Second should be second most recent"
763 );
764 })
765 }
766}