|
@@ -35,6 +35,8 @@ export interface InboxThread {
|
|
|
lastMessage: string
|
|
lastMessage: string
|
|
|
lastDirection: string
|
|
lastDirection: string
|
|
|
lastAt: string
|
|
lastAt: string
|
|
|
|
|
+ /** Inbound messages nobody has opened yet. */
|
|
|
|
|
+ unread: number
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
export interface InboxPage {
|
|
export interface InboxPage {
|
|
@@ -62,13 +64,22 @@ interface ThreadRow {
|
|
|
createdAt: Date
|
|
createdAt: Date
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+interface UnreadRow {
|
|
|
|
|
+ identity: string
|
|
|
|
|
+ unread: number
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
/** Empty bodies are common on media messages, so say what arrived instead. */
|
|
/** Empty bodies are common on media messages, so say what arrived instead. */
|
|
|
function preview(body: string | null, mediaType?: string | null): string {
|
|
function preview(body: string | null, mediaType?: string | null): string {
|
|
|
if (body?.trim()) return body
|
|
if (body?.trim()) return body
|
|
|
return mediaType ? `[${mediaType}]` : ''
|
|
return mediaType ? `[${mediaType}]` : ''
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
-function toThread(channel: MessagingChannel, row: ThreadRow): InboxThread {
|
|
|
|
|
|
|
+function toThread(
|
|
|
|
|
+ channel: MessagingChannel,
|
|
|
|
|
+ row: ThreadRow,
|
|
|
|
|
+ unread: Map<string, number>
|
|
|
|
|
+): InboxThread {
|
|
|
const contact = row.contact ?? ''
|
|
const contact = row.contact ?? ''
|
|
|
return {
|
|
return {
|
|
|
key: `${channel}:${row.identity}`,
|
|
key: `${channel}:${row.identity}`,
|
|
@@ -79,9 +90,14 @@ function toThread(channel: MessagingChannel, row: ThreadRow): InboxThread {
|
|
|
lastMessage: preview(row.body, row.mediaType),
|
|
lastMessage: preview(row.body, row.mediaType),
|
|
|
lastDirection: row.direction,
|
|
lastDirection: row.direction,
|
|
|
lastAt: row.createdAt.toISOString(),
|
|
lastAt: row.createdAt.toISOString(),
|
|
|
|
|
+ unread: unread.get(row.identity) ?? 0,
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+function byIdentity(rows: UnreadRow[]): Map<string, number> {
|
|
|
|
|
+ return new Map(rows.map((row) => [row.identity, row.unread]))
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
export interface InboxQuery {
|
|
export interface InboxQuery {
|
|
|
/** ISO timestamp of the last row already shown. */
|
|
/** ISO timestamp of the last row already shown. */
|
|
|
cursor?: string | null
|
|
cursor?: string | null
|
|
@@ -121,8 +137,13 @@ export async function getInboxThreads(query: InboxQuery = {}) {
|
|
|
// (organizationId, customerId, createdAt DESC) indexes already order for.
|
|
// (organizationId, customerId, createdAt DESC) indexes already order for.
|
|
|
// History outlives configuration, so this does not check `channels`: a
|
|
// History outlives configuration, so this does not check `channels`: a
|
|
|
// workshop that switched a provider off can still read what it sent.
|
|
// workshop that switched a provider off can still read what it sent.
|
|
|
- const [smsRows, telegramRows, whatsappRows] = await Promise.all([
|
|
|
|
|
- db.$queryRaw<ThreadRow[]>`
|
|
|
|
|
|
|
+ //
|
|
|
|
|
+ // Unread is counted per thread in its own small query rather than
|
|
|
|
|
+ // folded into the DISTINCT ON, which only ever sees one row per thread.
|
|
|
|
|
+ // Unread rows are few, so the whole workshop is counted at once.
|
|
|
|
|
+ const [smsRows, telegramRows, whatsappRows, smsUnread, telegramUnread, whatsappUnread] =
|
|
|
|
|
+ await Promise.all([
|
|
|
|
|
+ db.$queryRaw<ThreadRow[]>`
|
|
|
SELECT * FROM (
|
|
SELECT * FROM (
|
|
|
SELECT DISTINCT ON (m."customerId")
|
|
SELECT DISTINCT ON (m."customerId")
|
|
|
m."customerId" AS identity,
|
|
m."customerId" AS identity,
|
|
@@ -145,7 +166,7 @@ export async function getInboxThreads(query: InboxQuery = {}) {
|
|
|
ORDER BY t."createdAt" DESC
|
|
ORDER BY t."createdAt" DESC
|
|
|
LIMIT ${limit}
|
|
LIMIT ${limit}
|
|
|
`,
|
|
`,
|
|
|
- db.$queryRaw<ThreadRow[]>`
|
|
|
|
|
|
|
+ db.$queryRaw<ThreadRow[]>`
|
|
|
SELECT * FROM (
|
|
SELECT * FROM (
|
|
|
SELECT DISTINCT ON (m."customerId")
|
|
SELECT DISTINCT ON (m."customerId")
|
|
|
m."customerId" AS identity,
|
|
m."customerId" AS identity,
|
|
@@ -168,10 +189,10 @@ export async function getInboxThreads(query: InboxQuery = {}) {
|
|
|
ORDER BY t."createdAt" DESC
|
|
ORDER BY t."createdAt" DESC
|
|
|
LIMIT ${limit}
|
|
LIMIT ${limit}
|
|
|
`,
|
|
`,
|
|
|
- // WhatsApp threads can belong to a number we never matched to a
|
|
|
|
|
- // customer, so they group by customer when there is one and by the
|
|
|
|
|
- // number itself when there is not.
|
|
|
|
|
- db.$queryRaw<ThreadRow[]>`
|
|
|
|
|
|
|
+ // WhatsApp threads can belong to a number we never matched to a
|
|
|
|
|
+ // customer, so they group by customer when there is one and by the
|
|
|
|
|
+ // number itself when there is not.
|
|
|
|
|
+ db.$queryRaw<ThreadRow[]>`
|
|
|
SELECT * FROM (
|
|
SELECT * FROM (
|
|
|
SELECT DISTINCT ON (COALESCE(m."customerId", CASE WHEN m."direction" = 'inbound' THEN m."fromNumber" ELSE m."toNumber" END))
|
|
SELECT DISTINCT ON (COALESCE(m."customerId", CASE WHEN m."direction" = 'inbound' THEN m."fromNumber" ELSE m."toNumber" END))
|
|
|
COALESCE(m."customerId", CASE WHEN m."direction" = 'inbound' THEN m."fromNumber" ELSE m."toNumber" END) AS identity,
|
|
COALESCE(m."customerId", CASE WHEN m."direction" = 'inbound' THEN m."fromNumber" ELSE m."toNumber" END) AS identity,
|
|
@@ -195,13 +216,34 @@ export async function getInboxThreads(query: InboxQuery = {}) {
|
|
|
ORDER BY t."createdAt" DESC
|
|
ORDER BY t."createdAt" DESC
|
|
|
LIMIT ${limit}
|
|
LIMIT ${limit}
|
|
|
`,
|
|
`,
|
|
|
- ])
|
|
|
|
|
|
|
+ db.$queryRaw<UnreadRow[]>`
|
|
|
|
|
+ SELECT "customerId" AS identity, COUNT(*)::int AS unread
|
|
|
|
|
+ FROM "sms_messages"
|
|
|
|
|
+ WHERE "organizationId" = ${organizationId}
|
|
|
|
|
+ AND "direction" = 'inbound' AND "readAt" IS NULL AND "customerId" IS NOT NULL
|
|
|
|
|
+ GROUP BY "customerId"
|
|
|
|
|
+ `,
|
|
|
|
|
+ db.$queryRaw<UnreadRow[]>`
|
|
|
|
|
+ SELECT "customerId" AS identity, COUNT(*)::int AS unread
|
|
|
|
|
+ FROM "telegram_messages"
|
|
|
|
|
+ WHERE "organizationId" = ${organizationId}
|
|
|
|
|
+ AND "direction" = 'inbound' AND "readAt" IS NULL AND "customerId" IS NOT NULL
|
|
|
|
|
+ GROUP BY "customerId"
|
|
|
|
|
+ `,
|
|
|
|
|
+ db.$queryRaw<UnreadRow[]>`
|
|
|
|
|
+ SELECT COALESCE("customerId", "fromNumber") AS identity, COUNT(*)::int AS unread
|
|
|
|
|
+ FROM "whatsapp_messages"
|
|
|
|
|
+ WHERE "organizationId" = ${organizationId}
|
|
|
|
|
+ AND "direction" = 'inbound' AND "readAt" IS NULL
|
|
|
|
|
+ GROUP BY COALESCE("customerId", "fromNumber")
|
|
|
|
|
+ `,
|
|
|
|
|
+ ])
|
|
|
|
|
|
|
|
const { threads, nextCursor } = mergeChannelPages(
|
|
const { threads, nextCursor } = mergeChannelPages(
|
|
|
[
|
|
[
|
|
|
- smsRows.map((row) => toThread('sms', row)),
|
|
|
|
|
- telegramRows.map((row) => toThread('telegram', row)),
|
|
|
|
|
- whatsappRows.map((row) => toThread('whatsapp', row)),
|
|
|
|
|
|
|
+ smsRows.map((row) => toThread('sms', row, byIdentity(smsUnread))),
|
|
|
|
|
+ telegramRows.map((row) => toThread('telegram', row, byIdentity(telegramUnread))),
|
|
|
|
|
+ whatsappRows.map((row) => toThread('whatsapp', row, byIdentity(whatsappUnread))),
|
|
|
],
|
|
],
|
|
|
limit
|
|
limit
|
|
|
)
|
|
)
|
|
@@ -215,3 +257,112 @@ export async function getInboxThreads(query: InboxQuery = {}) {
|
|
|
}
|
|
}
|
|
|
)
|
|
)
|
|
|
}
|
|
}
|
|
|
|
|
+
|
|
|
|
|
+/**
|
|
|
|
|
+ * Marks everything inbound in one thread as read.
|
|
|
|
|
+ *
|
|
|
|
|
+ * Called when the thread is opened. Reading a conversation is what the
|
|
|
|
|
+ * customers permission already grants, so no extra permission is asked for.
|
|
|
|
|
+ * Returns how many rows changed, so the caller can skip a refresh when it was
|
|
|
|
|
+ * nothing.
|
|
|
|
|
+ */
|
|
|
|
|
+export async function markThreadRead(thread: {
|
|
|
|
|
+ channel: MessagingChannel
|
|
|
|
|
+ customerId: string | null
|
|
|
|
|
+ /** The phone number a WhatsApp thread without a customer is filed under. */
|
|
|
|
|
+ contact: string
|
|
|
|
|
+}) {
|
|
|
|
|
+ return withAuth(
|
|
|
|
|
+ async ({ organizationId }): Promise<{ marked: number }> => {
|
|
|
|
|
+ const now = new Date()
|
|
|
|
|
+ const unread = { organizationId, direction: 'inbound', readAt: null }
|
|
|
|
|
+
|
|
|
|
|
+ if (thread.channel === 'whatsapp') {
|
|
|
|
|
+ const result = await db.whatsappMessage.updateMany({
|
|
|
|
|
+ where: thread.customerId
|
|
|
|
|
+ ? { ...unread, customerId: thread.customerId }
|
|
|
|
|
+ : { ...unread, customerId: null, fromNumber: thread.contact },
|
|
|
|
|
+ data: { readAt: now },
|
|
|
|
|
+ })
|
|
|
|
|
+ return { marked: result.count }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ if (!thread.customerId) return { marked: 0 }
|
|
|
|
|
+
|
|
|
|
|
+ const result =
|
|
|
|
|
+ thread.channel === 'sms'
|
|
|
|
|
+ ? await db.smsMessage.updateMany({
|
|
|
|
|
+ where: { ...unread, customerId: thread.customerId },
|
|
|
|
|
+ data: { readAt: now },
|
|
|
|
|
+ })
|
|
|
|
|
+ : await db.telegramMessage.updateMany({
|
|
|
|
|
+ where: { ...unread, customerId: thread.customerId },
|
|
|
|
|
+ data: { readAt: now },
|
|
|
|
|
+ })
|
|
|
|
|
+ return { marked: result.count }
|
|
|
|
|
+ },
|
|
|
|
|
+ {
|
|
|
|
|
+ requiredPermissions: [
|
|
|
|
|
+ { action: PermissionAction.READ, subject: PermissionSubject.CUSTOMERS },
|
|
|
|
|
+ ],
|
|
|
|
|
+ }
|
|
|
|
|
+ )
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+/**
|
|
|
|
|
+ * Puts a thread back in the unread pile.
|
|
|
|
|
+ *
|
|
|
|
|
+ * Only the newest inbound message is unstamped: one waiting message is what
|
|
|
|
|
+ * "come back to this" means, and it keeps the pill honest about how much is
|
|
|
|
|
+ * actually new. A thread the customer has never written to has nothing to
|
|
|
|
|
+ * mark, and says so with a zero.
|
|
|
|
|
+ */
|
|
|
|
|
+export async function markThreadUnread(thread: {
|
|
|
|
|
+ channel: MessagingChannel
|
|
|
|
|
+ customerId: string | null
|
|
|
|
|
+ contact: string
|
|
|
|
|
+}) {
|
|
|
|
|
+ return withAuth(
|
|
|
|
|
+ async ({ organizationId }): Promise<{ marked: number }> => {
|
|
|
|
|
+ const inbound = { organizationId, direction: 'inbound' }
|
|
|
|
|
+ const newest = { orderBy: { createdAt: 'desc' as const }, select: { id: true } }
|
|
|
|
|
+
|
|
|
|
|
+ if (thread.channel === 'whatsapp') {
|
|
|
|
|
+ const latest = await db.whatsappMessage.findFirst({
|
|
|
|
|
+ where: thread.customerId
|
|
|
|
|
+ ? { ...inbound, customerId: thread.customerId }
|
|
|
|
|
+ : { ...inbound, customerId: null, fromNumber: thread.contact },
|
|
|
|
|
+ ...newest,
|
|
|
|
|
+ })
|
|
|
|
|
+ if (!latest) return { marked: 0 }
|
|
|
|
|
+ await db.whatsappMessage.update({ where: { id: latest.id }, data: { readAt: null } })
|
|
|
|
|
+ return { marked: 1 }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ if (!thread.customerId) return { marked: 0 }
|
|
|
|
|
+
|
|
|
|
|
+ if (thread.channel === 'sms') {
|
|
|
|
|
+ const latest = await db.smsMessage.findFirst({
|
|
|
|
|
+ where: { ...inbound, customerId: thread.customerId },
|
|
|
|
|
+ ...newest,
|
|
|
|
|
+ })
|
|
|
|
|
+ if (!latest) return { marked: 0 }
|
|
|
|
|
+ await db.smsMessage.update({ where: { id: latest.id }, data: { readAt: null } })
|
|
|
|
|
+ return { marked: 1 }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ const latest = await db.telegramMessage.findFirst({
|
|
|
|
|
+ where: { ...inbound, customerId: thread.customerId },
|
|
|
|
|
+ ...newest,
|
|
|
|
|
+ })
|
|
|
|
|
+ if (!latest) return { marked: 0 }
|
|
|
|
|
+ await db.telegramMessage.update({ where: { id: latest.id }, data: { readAt: null } })
|
|
|
|
|
+ return { marked: 1 }
|
|
|
|
|
+ },
|
|
|
|
|
+ {
|
|
|
|
|
+ requiredPermissions: [
|
|
|
|
|
+ { action: PermissionAction.READ, subject: PermissionSubject.CUSTOMERS },
|
|
|
|
|
+ ],
|
|
|
|
|
+ }
|
|
|
|
|
+ )
|
|
|
|
|
+}
|