From ac3af93df58d88d2b1fbeeff36aa22f58f619ef3 Mon Sep 17 00:00:00 2001 From: viettranx Date: Thu, 2 Apr 2026 08:23:01 +0700 Subject: [PATCH] feat(contacts): implement thread_id persistence in PG and SQLite stores - Update UpsertContact to handle threadID and threadType - Strip username from sender_id compound key - Implement in both PostgreSQL and SQLite backends --- internal/store/pg/channel_contacts.go | 29 ++++++++++++++------------ internal/store/sqlitestore/contacts.go | 22 ++++++++++--------- 2 files changed, 28 insertions(+), 23 deletions(-) diff --git a/internal/store/pg/channel_contacts.go b/internal/store/pg/channel_contacts.go index 073301d1..807123fd 100644 --- a/internal/store/pg/channel_contacts.go +++ b/internal/store/pg/channel_contacts.go @@ -23,23 +23,24 @@ func NewPGContactStore(db *sql.DB) *PGContactStore { return &PGContactStore{db: db, resolveCache: newContactResolveCache()} } -func (s *PGContactStore) UpsertContact(ctx context.Context, channelType, channelInstance, senderID, userID, displayName, username, peerKind, contactType string) error { +func (s *PGContactStore) UpsertContact(ctx context.Context, channelType, channelInstance, senderID, userID, displayName, username, peerKind, contactType, threadID, threadType string) error { tenantID := store.TenantIDFromContext(ctx) if tenantID == uuid.Nil { tenantID = store.MasterTenantID } _, err := s.db.ExecContext(ctx, ` - INSERT INTO channel_contacts (channel_type, channel_instance, sender_id, user_id, display_name, username, peer_kind, contact_type, tenant_id) - VALUES ($1, NULLIF($2,''), $3, NULLIF($4,''), NULLIF($5,''), NULLIF($6,''), NULLIF($7,''), $8, $9) - ON CONFLICT (tenant_id, channel_type, sender_id) DO UPDATE SET + INSERT INTO channel_contacts (channel_type, channel_instance, sender_id, user_id, display_name, username, peer_kind, contact_type, thread_id, thread_type, tenant_id) + VALUES ($1, NULLIF($2,''), $3, NULLIF($4,''), NULLIF($5,''), NULLIF($6,''), NULLIF($7,''), $8, NULLIF($9,''), NULLIF($10,''), $11) + ON CONFLICT (tenant_id, channel_type, sender_id, COALESCE(thread_id, '')) DO UPDATE SET display_name = COALESCE(NULLIF($5,''), channel_contacts.display_name), username = COALESCE(NULLIF($6,''), channel_contacts.username), user_id = COALESCE(NULLIF($4,''), channel_contacts.user_id), channel_instance = COALESCE(NULLIF($2,''), channel_contacts.channel_instance), peer_kind = COALESCE(NULLIF($7,''), channel_contacts.peer_kind), contact_type = $8, + thread_type = COALESCE(NULLIF($10,''), channel_contacts.thread_type), last_seen_at = NOW()`, - channelType, channelInstance, senderID, userID, displayName, username, peerKind, contactType, tenantID, + channelType, channelInstance, senderID, userID, displayName, username, peerKind, contactType, threadID, threadType, tenantID, ) return err } @@ -95,7 +96,7 @@ func (s *PGContactStore) ListContacts(ctx context.Context, opts store.ContactLis where, args, argIdx := contactWhereClause(ctx, opts) query := `SELECT id, channel_type, channel_instance, sender_id, user_id, - display_name, username, avatar_url, peer_kind, contact_type, merged_id, + display_name, username, avatar_url, peer_kind, contact_type, thread_id, thread_type, merged_id, first_seen_at, last_seen_at FROM channel_contacts` + where + " ORDER BY last_seen_at DESC" @@ -123,7 +124,7 @@ func (s *PGContactStore) ListContacts(ctx context.Context, opts store.ContactLis var c store.ChannelContact if err := rows.Scan( &c.ID, &c.ChannelType, &c.ChannelInstance, &c.SenderID, &c.UserID, - &c.DisplayName, &c.Username, &c.AvatarURL, &c.PeerKind, &c.ContactType, &c.MergedID, + &c.DisplayName, &c.Username, &c.AvatarURL, &c.PeerKind, &c.ContactType, &c.ThreadID, &c.ThreadType, &c.MergedID, &c.FirstSeenAt, &c.LastSeenAt, ); err != nil { return nil, err @@ -154,7 +155,7 @@ func (s *PGContactStore) GetContactsBySenderIDs(ctx context.Context, senderIDs [ query := fmt.Sprintf(`SELECT DISTINCT ON (sender_id) id, channel_type, channel_instance, sender_id, user_id, - display_name, username, avatar_url, peer_kind, contact_type, merged_id, + display_name, username, avatar_url, peer_kind, contact_type, thread_id, thread_type, merged_id, first_seen_at, last_seen_at FROM channel_contacts WHERE sender_id IN (%s) @@ -171,7 +172,7 @@ func (s *PGContactStore) GetContactsBySenderIDs(ctx context.Context, senderIDs [ var c store.ChannelContact if err := rows.Scan( &c.ID, &c.ChannelType, &c.ChannelInstance, &c.SenderID, &c.UserID, - &c.DisplayName, &c.Username, &c.AvatarURL, &c.PeerKind, &c.ContactType, &c.MergedID, + &c.DisplayName, &c.Username, &c.AvatarURL, &c.PeerKind, &c.ContactType, &c.ThreadID, &c.ThreadType, &c.MergedID, &c.FirstSeenAt, &c.LastSeenAt, ); err != nil { return nil, err @@ -185,13 +186,15 @@ func (s *PGContactStore) GetContactByID(ctx context.Context, id uuid.UUID) (*sto tid := store.TenantIDFromContext(ctx) row := s.db.QueryRowContext(ctx, `SELECT id, channel_type, channel_instance, sender_id, user_id, - display_name, username, avatar_url, peer_kind, merged_id, + display_name, username, avatar_url, peer_kind, contact_type, + thread_id, thread_type, merged_id, first_seen_at, last_seen_at FROM channel_contacts WHERE id = $1 AND tenant_id = $2`, id, tid) var c store.ChannelContact if err := row.Scan( &c.ID, &c.ChannelType, &c.ChannelInstance, &c.SenderID, &c.UserID, - &c.DisplayName, &c.Username, &c.AvatarURL, &c.PeerKind, &c.MergedID, + &c.DisplayName, &c.Username, &c.AvatarURL, &c.PeerKind, &c.ContactType, + &c.ThreadID, &c.ThreadType, &c.MergedID, &c.FirstSeenAt, &c.LastSeenAt, ); err != nil { return nil, err @@ -286,7 +289,7 @@ func (s *PGContactStore) GetContactsByMergedID(ctx context.Context, mergedID uui tid := store.TenantIDFromContext(ctx) q := `SELECT id, channel_type, channel_instance, sender_id, user_id, - display_name, username, avatar_url, peer_kind, contact_type, merged_id, + display_name, username, avatar_url, peer_kind, contact_type, thread_id, thread_type, merged_id, first_seen_at, last_seen_at FROM channel_contacts WHERE merged_id = $1 AND tenant_id = $2 ORDER BY last_seen_at DESC` @@ -302,7 +305,7 @@ func (s *PGContactStore) GetContactsByMergedID(ctx context.Context, mergedID uui var c store.ChannelContact if err := rows.Scan( &c.ID, &c.ChannelType, &c.ChannelInstance, &c.SenderID, &c.UserID, - &c.DisplayName, &c.Username, &c.AvatarURL, &c.PeerKind, &c.ContactType, &c.MergedID, + &c.DisplayName, &c.Username, &c.AvatarURL, &c.PeerKind, &c.ContactType, &c.ThreadID, &c.ThreadType, &c.MergedID, &c.FirstSeenAt, &c.LastSeenAt, ); err != nil { return nil, err diff --git a/internal/store/sqlitestore/contacts.go b/internal/store/sqlitestore/contacts.go index 93ed8077..afc6c334 100644 --- a/internal/store/sqlitestore/contacts.go +++ b/internal/store/sqlitestore/contacts.go @@ -23,21 +23,22 @@ func NewSQLiteContactStore(db *sql.DB) *SQLiteContactStore { return &SQLiteContactStore{db: db} } -func (s *SQLiteContactStore) UpsertContact(ctx context.Context, channelType, channelInstance, senderID, userID, displayName, username, peerKind, contactType string) error { +func (s *SQLiteContactStore) UpsertContact(ctx context.Context, channelType, channelInstance, senderID, userID, displayName, username, peerKind, contactType, threadID, threadType string) error { tenantID := store.TenantIDFromContext(ctx) if tenantID == uuid.Nil { tenantID = store.MasterTenantID } _, err := s.db.ExecContext(ctx, ` - INSERT INTO channel_contacts (channel_type, channel_instance, sender_id, user_id, display_name, username, peer_kind, contact_type, tenant_id) - VALUES (?, NULLIF(?,?), ?, NULLIF(?,?), NULLIF(?,?), NULLIF(?,?), NULLIF(?,?), ?, ?) - ON CONFLICT (tenant_id, channel_type, sender_id) DO UPDATE SET + INSERT INTO channel_contacts (channel_type, channel_instance, sender_id, user_id, display_name, username, peer_kind, contact_type, thread_id, thread_type, tenant_id) + VALUES (?, NULLIF(?,?), ?, NULLIF(?,?), NULLIF(?,?), NULLIF(?,?), NULLIF(?,?), ?, NULLIF(?,?), NULLIF(?,?), ?) + ON CONFLICT (tenant_id, channel_type, sender_id, COALESCE(thread_id, '')) DO UPDATE SET display_name = COALESCE(NULLIF(excluded.display_name,''), channel_contacts.display_name), username = COALESCE(NULLIF(excluded.username,''), channel_contacts.username), user_id = COALESCE(NULLIF(excluded.user_id,''), channel_contacts.user_id), channel_instance = COALESCE(NULLIF(excluded.channel_instance,''), channel_contacts.channel_instance), peer_kind = COALESCE(NULLIF(excluded.peer_kind,''), channel_contacts.peer_kind), contact_type = excluded.contact_type, + thread_type = COALESCE(NULLIF(excluded.thread_type,''), channel_contacts.thread_type), last_seen_at = CURRENT_TIMESTAMP`, channelType, channelInstance, "", @@ -47,6 +48,8 @@ func (s *SQLiteContactStore) UpsertContact(ctx context.Context, channelType, cha username, "", peerKind, "", contactType, + threadID, "", + threadType, "", tenantID, ) return err @@ -91,14 +94,14 @@ func contactWhereSQLite(ctx context.Context, opts store.ContactListOpts) (string } const contactSelectCols = `id, channel_type, channel_instance, sender_id, user_id, - display_name, username, avatar_url, peer_kind, contact_type, merged_id, + display_name, username, avatar_url, peer_kind, contact_type, thread_id, thread_type, merged_id, first_seen_at, last_seen_at` func scanContact(rows *sql.Rows) (store.ChannelContact, error) { var c store.ChannelContact err := rows.Scan( &c.ID, &c.ChannelType, &c.ChannelInstance, &c.SenderID, &c.UserID, - &c.DisplayName, &c.Username, &c.AvatarURL, &c.PeerKind, &c.ContactType, &c.MergedID, + &c.DisplayName, &c.Username, &c.AvatarURL, &c.PeerKind, &c.ContactType, &c.ThreadID, &c.ThreadType, &c.MergedID, &c.FirstSeenAt, &c.LastSeenAt, ) return c, err @@ -154,9 +157,7 @@ func (s *SQLiteContactStore) GetContactsBySenderIDs(ctx context.Context, senderI } // SQLite has no DISTINCT ON; emulate with GROUP BY + MAX rowid trick via subquery - query := `SELECT id, channel_type, channel_instance, sender_id, user_id, - display_name, username, avatar_url, peer_kind, contact_type, merged_id, - first_seen_at, last_seen_at + query := `SELECT ` + contactSelectCols + ` FROM channel_contacts WHERE sender_id IN (` + strings.Join(placeholders, ",") + `) GROUP BY sender_id @@ -187,7 +188,8 @@ func (s *SQLiteContactStore) GetContactByID(ctx context.Context, id uuid.UUID) ( var c store.ChannelContact if err := row.Scan( &c.ID, &c.ChannelType, &c.ChannelInstance, &c.SenderID, &c.UserID, - &c.DisplayName, &c.Username, &c.AvatarURL, &c.PeerKind, &c.ContactType, &c.MergedID, + &c.DisplayName, &c.Username, &c.AvatarURL, &c.PeerKind, &c.ContactType, + &c.ThreadID, &c.ThreadType, &c.MergedID, &c.FirstSeenAt, &c.LastSeenAt, ); err != nil { return nil, err