diff --git a/hasura/migrations/default/1889000001600_drop_chat_notifications/down.sql b/hasura/migrations/default/1889000001600_drop_chat_notifications/down.sql new file mode 100644 index 00000000..bc4f69a8 --- /dev/null +++ b/hasura/migrations/default/1889000001600_drop_chat_notifications/down.sql @@ -0,0 +1,4 @@ +-- The rows are gone for good: chat pushes straight from the conversation. +CREATE INDEX IF NOT EXISTS notifications_message_id_idx + ON public.notifications ((data->>'messageId')) + WHERE data->>'messageId' IS NOT NULL; diff --git a/hasura/migrations/default/1889000001600_drop_chat_notifications/up.sql b/hasura/migrations/default/1889000001600_drop_chat_notifications/up.sql new file mode 100644 index 00000000..0056fd56 --- /dev/null +++ b/hasura/migrations/default/1889000001600_drop_chat_notifications/up.sql @@ -0,0 +1,8 @@ +DELETE FROM public.notifications + WHERE type IN ('ChatMessage', 'MatchChatMessage'); + +DELETE FROM public.notification_preferences + WHERE channel = 'in_app' + AND key IN ('ChatMessage', 'MatchChatMessage'); + +DROP INDEX IF EXISTS public.notifications_message_id_idx; diff --git a/hasura/triggers/player_blocks.sql b/hasura/triggers/player_blocks.sql index 4df62068..74d51ca2 100644 --- a/hasura/triggers/player_blocks.sql +++ b/hasura/triggers/player_blocks.sql @@ -126,29 +126,6 @@ BEGIN WHERE dc.room_id = LEAST(_a, _b)::text || ':' || GREATEST(_a, _b)::text AND dc.steam_id = _a; - -- Chat hides what the blocked player said from the blocker, so their bell - -- previews go the way NotificationsService.retractChatMessage takes them: - -- blanked as well, since a recipient can restore their own deleted rows, - -- and a push still waiting in its bundling window skips a deleted row. - -- A moderator still sees group rooms whole (ChatService.blockExemptRoles), - -- so only their DM rows go. - UPDATE public.notifications n - SET deleted_at = COALESCE(n.deleted_at, now()), - message = '' - WHERE n.steam_id = _a - AND n.type IN ('ChatMessage', 'MatchChatMessage') - AND n.data->>'senderSteamId' = _b::text - AND (n.deleted_at IS NULL OR n.message <> '') - AND ( - n.entity_id LIKE 'direct:%' - OR NOT EXISTS ( - SELECT 1 - FROM public.players p - WHERE p.steam_id = _a - AND public.is_role_below('moderator', p.role::text) - ) - ); - RETURN NULL; END; $$; diff --git a/src/chat/chat.service.spec.ts b/src/chat/chat.service.spec.ts index d8546573..4e204627 100644 --- a/src/chat/chat.service.spec.ts +++ b/src/chat/chat.service.spec.ts @@ -204,13 +204,10 @@ describe("ChatService direct messages", () => { }), }; - const notifications = { - notifyPlayers: jest.fn(), - collapseOlderUnread: jest.fn(), - markConversationRead: jest.fn(), + const push = { + sendChatMessage: jest.fn(), retractChatMessage: jest.fn().mockResolvedValue(undefined), - retractChatMessageFromBlocked: jest.fn().mockResolvedValue(undefined), - updateChatMessagePreview: jest.fn().mockResolvedValue(undefined), + editChatMessage: jest.fn().mockResolvedValue(undefined), }; const client = (steamId: string) => @@ -425,9 +422,8 @@ describe("ChatService direct messages", () => { redis.eval.mockImplementation(async (script: string) => script.includes("INCR") ? 1 : [1, 1], ); - notifications.retractChatMessage.mockResolvedValue(undefined); - notifications.retractChatMessageFromBlocked.mockResolvedValue(undefined); - notifications.updateChatMessagePreview.mockResolvedValue(undefined); + push.retractChatMessage.mockResolvedValue(undefined); + push.editChatMessage.mockResolvedValue(undefined); directMessage = undefined; acceptedFriendships = [[ME, FRIEND]]; myMatches = ["m-1"]; @@ -458,7 +454,7 @@ describe("ChatService direct messages", () => { hasuraService as any, postgres as any, { getConnection: () => redis } as any, - notifications as any, + push as any, playerBlocks as any, ); }); @@ -765,7 +761,7 @@ describe("ChatService direct messages", () => { expect(ChatService.notificationTypeFor(type)).toBe("ChatMessage"); }); - it("notifies and collapses a team room line under one type", async () => { + it("pushes a team room line as match chat", async () => { await service.sendMessageToChat( ChatLobbyType.MatchTeam, "m-1:l-1", @@ -779,30 +775,14 @@ describe("ChatService direct messages", () => { await new Promise((resolve) => setImmediate(resolve)); } - expect(notifications.notifyPlayers).toHaveBeenCalledWith( - "MatchChatMessage", + expect(push.sendChatMessage).toHaveBeenCalledWith( + [FRIEND], expect.objectContaining({ - entity_id: "match_team:m-1:l-1", - steamIds: [FRIEND], + type: "MatchChatMessage", + entityId: "match_team:m-1:l-1", + threadKey: "chat:match_team:m-1:l-1", }), ); - expect(notifications.collapseOlderUnread).toHaveBeenCalledWith( - "MatchChatMessage", - "match_team:m-1:l-1", - [FRIEND], - ); - }); - - it("clears a team room's badge under the type it was sent with", async () => { - await service.markThreadRead(ChatLobbyType.MatchTeam, "m-1:l-1", { - steam_id: ME, - } as any); - - expect(notifications.markConversationRead).toHaveBeenCalledWith( - "MatchChatMessage", - "match_team:m-1:l-1", - ME, - ); }); }); @@ -1206,7 +1186,7 @@ describe("ChatService direct messages", () => { expect(redis.hset).not.toHaveBeenCalled(); expect(redis.publish).not.toHaveBeenCalled(); - expect(notifications.notifyPlayers).not.toHaveBeenCalled(); + expect(push.sendChatMessage).not.toHaveBeenCalled(); }, ); @@ -1479,7 +1459,7 @@ describe("ChatService direct messages", () => { ).rejects.toThrow("database down"); expect(redis.hdel).not.toHaveBeenCalled(); - expect(notifications.retractChatMessage).not.toHaveBeenCalled(); + expect(push.retractChatMessage).not.toHaveBeenCalled(); }); it("records no author for a steam id stored as a number", async () => { @@ -1507,12 +1487,12 @@ describe("ChatService direct messages", () => { moderator(), ); - expect(notifications.retractChatMessage).toHaveBeenCalledWith(MESSAGE_ID); + expect(push.retractChatMessage).toHaveBeenCalledWith(MESSAGE_ID); }); it("still deletes when the retraction fails", async () => { store(ChatLobbyType.Match, "m-1"); - notifications.retractChatMessage.mockRejectedValue(new Error("nope")); + push.retractChatMessage.mockRejectedValue(new Error("nope")); await expect( service.deleteMessage( @@ -1623,58 +1603,14 @@ describe("ChatService direct messages", () => { expect(redis.hget).not.toHaveBeenCalled(); }); - describe("while its notifications are still being written", () => { - const sayInTournament = async () => { - tournament.roster = [ME, FRIEND]; - redis.hget.mockResolvedValue( - JSON.stringify({ user: { steam_id: ME } }), - ); - role = "user"; - - const result = await service.sendMessageToChat( - ChatLobbyType.Tournament, - "t-1", - { steam_id: ME, name: "Someone", role: "user" } as any, - "hi", - ); - - await flush(); - await flush(); - - return result.accepted ? result.messageId : undefined; - }; - - it("retracts them once written if the message was deleted meanwhile", async () => { - audited = true; - - const messageId = await sayInTournament(); - - expect(notifications.notifyPlayers).toHaveBeenCalled(); - expect( - queries.find(({ sql }) => - sql.includes("SELECT 1 FROM public.chat_message_deletions"), - )?.bindings, - ).toEqual([messageId]); - expect(notifications.retractChatMessage).toHaveBeenCalledWith( - messageId, - ); - }); - - it("leaves them alone when nothing deleted the message", async () => { - await sayInTournament(); - - expect(notifications.notifyPlayers).toHaveBeenCalled(); - expect(notifications.retractChatMessage).not.toHaveBeenCalled(); - }); - }); - - it("stamps each chat notification with the message it announces", async () => { + it("pushes a direct message straight to the other party", async () => { redis.hget.mockResolvedValue(JSON.stringify({ user: { steam_id: ME } })); role = "user"; + const room = directRoomId(ME, FRIEND); const result = await service.sendMessageToChat( ChatLobbyType.Direct, - directRoomId(ME, FRIEND), + room, { steam_id: ME, name: "Someone", role: "user" } as any, "hi", ); @@ -1682,14 +1618,18 @@ describe("ChatService direct messages", () => { await flush(); expect(result.accepted).toBe(true); - expect(notifications.notifyPlayers).toHaveBeenCalledWith( - "ChatMessage", - expect.objectContaining({ - data: expect.objectContaining({ - messageId: result.accepted ? result.messageId : undefined, - }), - }), - ); + expect(push.sendChatMessage).toHaveBeenCalledWith([FRIEND], { + messageId: result.accepted ? result.messageId : undefined, + type: "ChatMessage", + title: "Someone", + message: "hi", + entityId: `direct:${room}`, + threadKey: `chat:direct:${room}`, + threadLabel: "Someone", + icon: undefined, + senderSteamId: ME, + blockExemptRoles: [], + }); }); }); @@ -1837,24 +1777,22 @@ describe("ChatService direct messages", () => { expect(broadcasts("lobby:match:m-1:chat")).toEqual([]); }); - it("rewrites the bell's preview rather than notifying again", async () => { + it("rewrites a held push's text rather than notifying again", async () => { store(); await edit("fixed"); await flush(); - expect(notifications.updateChatMessagePreview).toHaveBeenCalledWith( + expect(push.editChatMessage).toHaveBeenCalledWith( MESSAGE_ID, "<b>fixed</b>", ); - expect(notifications.notifyPlayers).not.toHaveBeenCalled(); + expect(push.sendChatMessage).not.toHaveBeenCalled(); }); it("still edits when the preview cannot be rewritten", async () => { store(); - notifications.updateChatMessagePreview.mockRejectedValue( - new Error("database down"), - ); + push.editChatMessage.mockRejectedValue(new Error("database down")); await expect(edit()).resolves.toMatchObject({ edited: true }); expect(current().message).toBe("fixed"); @@ -2159,9 +2097,7 @@ describe("ChatService direct messages", () => { "web", ME, ]); - expect(notifications.retractChatMessage).toHaveBeenCalledWith( - MESSAGE_ID, - ); + expect(push.retractChatMessage).toHaveBeenCalledWith(MESSAGE_ID); }); it("lets a gagged author delete their own message", async () => { @@ -2269,10 +2205,7 @@ describe("ChatService direct messages", () => { }, }, ]); - expect(notifications.updateChatMessagePreview).toHaveBeenCalledWith( - MESSAGE_ID, - "fixed", - ); + expect(push.editChatMessage).toHaveBeenCalledWith(MESSAGE_ID, "fixed"); }); it("never audits the edit of a direct message", async () => { @@ -2373,9 +2306,7 @@ describe("ChatService direct messages", () => { data: { id: MESSAGE_ID }, }, ]); - expect(notifications.retractChatMessage).toHaveBeenCalledWith( - MESSAGE_ID, - ); + expect(push.retractChatMessage).toHaveBeenCalledWith(MESSAGE_ID); }); it("refuses to delete the other party's message", async () => { @@ -2388,131 +2319,6 @@ describe("ChatService direct messages", () => { expect(directMessage).toBeDefined(); }); }); - - describe("while its notifications are still being written", () => { - const say = async (type: ChatLobbyType, id: string) => { - redis.hget.mockImplementation(async (key: string, field: string) => - key.startsWith("chat:") - ? JSON.stringify({ user: { steam_id: ME } }) - : (stored[key]?.[field] ?? null), - ); - - const result = await service.sendMessageToChat(type, id, me(), "typo"); - - return result.accepted ? result.messageId : undefined; - }; - - it("gives the rows a room message's edited text", async () => { - tournament.roster = [ME, FRIEND]; - notifications.notifyPlayers.mockImplementationOnce(async () => { - const [key] = redis.hset.mock.calls.at(-1); - const [field, raw] = redis.hset.mock.calls.at(-1).slice(1); - stored[key] = { - [field]: JSON.stringify({ - ...JSON.parse(raw), - message: "fixed", - edited_at: new Date().toISOString(), - }), - }; - }); - - const messageId = await say(ChatLobbyType.Tournament, "t-1"); - await flush(); - await flush(); - - expect(notifications.updateChatMessagePreview).toHaveBeenCalledWith( - messageId, - "fixed", - ); - expect(notifications.retractChatMessage).not.toHaveBeenCalled(); - }); - - it("leaves the rows alone when the message is unchanged", async () => { - tournament.roster = [ME, FRIEND]; - redis.hset.mockImplementationOnce( - async (key: string, field: string, raw: string) => { - stored[key] = { [field]: raw }; - }, - ); - - await say(ChatLobbyType.Tournament, "t-1"); - await flush(); - await flush(); - - expect(notifications.notifyPlayers).toHaveBeenCalled(); - expect(notifications.updateChatMessagePreview).not.toHaveBeenCalled(); - expect(notifications.retractChatMessage).not.toHaveBeenCalled(); - }); - - it("retracts them when a direct message was deleted meanwhile", async () => { - const room = directRoomId(ME, FRIEND); - notifications.notifyPlayers.mockImplementationOnce(async () => { - const insert = queries.find(({ sql }) => - sql.includes("INSERT INTO public.direct_messages"), - ); - directMessage = { - id: insert.bindings[0], - roomId: room, - author: ME, - message: "typo", - open: true, - editedAt: null, - }; - - await expect( - service.deleteMessage( - ChatLobbyType.Direct, - room, - insert.bindings[0], - me(), - ), - ).resolves.toEqual({ deleted: true }); - }); - - const messageId = await say(ChatLobbyType.Direct, room); - await flush(); - await flush(); - - expect(notifications.retractChatMessage).toHaveBeenCalledTimes(2); - expect(notifications.retractChatMessage).toHaveBeenLastCalledWith( - messageId, - ); - expect( - notifications.retractChatMessage.mock.invocationCallOrder[1], - ).toBeGreaterThan( - notifications.collapseOlderUnread.mock.invocationCallOrder[0], - ); - }); - - it("gives them a direct message's edited text", async () => { - notifications.notifyPlayers.mockImplementationOnce(async () => { - const insert = queries.find(({ sql }) => - sql.includes("INSERT INTO public.direct_messages"), - ); - directMessage = { - id: insert.bindings[0], - roomId: insert.bindings[1], - author: ME, - message: "fixed", - open: true, - editedAt: new Date(), - }; - }); - - const messageId = await say( - ChatLobbyType.Direct, - directRoomId(ME, FRIEND), - ); - await flush(); - await flush(); - - expect(notifications.updateChatMessagePreview).toHaveBeenCalledWith( - messageId, - "fixed", - ); - expect(notifications.retractChatMessage).not.toHaveBeenCalled(); - }); - }); }); describe("reactions", () => { @@ -2852,7 +2658,7 @@ describe("ChatService direct messages", () => { await react(); await flush(); - expect(notifications.notifyPlayers).not.toHaveBeenCalled(); + expect(push.sendChatMessage).not.toHaveBeenCalled(); expect(rcon.connect).not.toHaveBeenCalled(); expect(broadcasts("lobby:match:m-1:chat")).toEqual([]); }); @@ -2910,7 +2716,7 @@ describe("ChatService direct messages", () => { await flush(); expect(broadcasts("lobby:match:m-1:deleted")).toHaveLength(1); - expect(notifications.retractChatMessage).toHaveBeenCalledWith(MESSAGE_ID); + expect(push.retractChatMessage).toHaveBeenCalledWith(MESSAGE_ID); }); it("puts each message's reactions in the room's history", async () => { @@ -3488,7 +3294,7 @@ describe("ChatService direct messages", () => { ), ).toBe(false); expect(redis.publish).not.toHaveBeenCalled(); - expect(notifications.notifyPlayers).not.toHaveBeenCalled(); + expect(push.sendChatMessage).not.toHaveBeenCalled(); }); it("keeps a direct message from a player its recipient has just blocked", async () => { @@ -3503,7 +3309,7 @@ describe("ChatService direct messages", () => { expect(published().map(({ steamId }) => steamId)).not.toContain(FRIEND); expect(playerBlocks.hasBlocked).toHaveBeenCalledWith(FRIEND, ME); - expect(notifications.notifyPlayers).not.toHaveBeenCalled(); + expect(push.sendChatMessage).not.toHaveBeenCalled(); }); it("keeps the socket's cleanup when the history's block lookup fails", async () => { @@ -3867,23 +3673,22 @@ describe("ChatService direct messages", () => { return result.accepted ? result.messageId : undefined; }; - it("writes none for a player who blocked the sender", async () => { + it("pushes nothing to a player who blocked the sender", async () => { blocks = [[ME, FRIEND]]; await sayInTournament(FRIEND); - expect(notifications.notifyPlayers).toHaveBeenCalledWith( - "ChatMessage", - expect.objectContaining({ steamIds: [STRANGER] }), - ); - expect(notifications.collapseOlderUnread).toHaveBeenCalledWith( - "ChatMessage", - "tournament:t-1", + expect(push.sendChatMessage).toHaveBeenCalledWith( [STRANGER], + expect.objectContaining({ + type: "ChatMessage", + senderSteamId: FRIEND, + blockExemptRoles: MODERATORS, + }), ); }); - it("writes none at all when everyone else blocked the sender", async () => { + it("pushes nothing at all when everyone else blocked the sender", async () => { blocks = [ [ME, FRIEND], [STRANGER, FRIEND], @@ -3896,32 +3701,18 @@ describe("ChatService direct messages", () => { [FRIEND], MODERATORS, ); - expect(notifications.notifyPlayers).not.toHaveBeenCalled(); + expect(push.sendChatMessage).not.toHaveBeenCalled(); expect(logger.warn).not.toHaveBeenCalled(); }); - it("catches up on a block that landed while the rows were being written", async () => { - const messageId = await sayInTournament(FRIEND); - - expect( - notifications.retractChatMessageFromBlocked, - ).toHaveBeenCalledWith(messageId, MODERATORS); - expect( - notifications.retractChatMessageFromBlocked.mock - .invocationCallOrder[0], - ).toBeGreaterThan( - notifications.notifyPlayers.mock.invocationCallOrder[0], - ); - }); - it("still notifies the player a block is aimed at", async () => { blocks = [[FRIEND, ME]]; await sayInTournament(FRIEND); - expect( - notifications.notifyPlayers.mock.calls[0][1].steamIds.sort(), - ).toEqual([ME, STRANGER].sort()); + expect(push.sendChatMessage.mock.calls[0][0].sort()).toEqual( + [ME, STRANGER].sort(), + ); }); }); diff --git a/src/chat/chat.service.ts b/src/chat/chat.service.ts index 60348079..1ea33234 100644 --- a/src/chat/chat.service.ts +++ b/src/chat/chat.service.ts @@ -14,6 +14,7 @@ import { } from "generated/schema"; import { isRoleAbove, rolesAtOrAbove } from "src/utilities/isRoleAbove"; import { NotificationsService } from "src/notifications/notifications.service"; +import { PushNotificationsService } from "src/notifications/push/push-notifications.service"; import { PostgresService } from "src/postgres/postgres.service"; import { PlayerBlocksService } from "src/player-blocks/player-blocks.service"; import { chatThreadKey } from "src/notifications/push/notification-delivery"; @@ -302,7 +303,7 @@ export class ChatService { private readonly hasuraService: HasuraService, private readonly postgres: PostgresService, private readonly redisManager: RedisManagerService, - private readonly notifications: NotificationsService, + private readonly pushNotifications: PushNotificationsService, private readonly playerBlocks: PlayerBlocksService, ) { this.redis = this.redisManager.getConnection(); @@ -1595,8 +1596,8 @@ export class ChatService { return null; } - // Never relayed to the game server, and no new notification: the bell rows - // already announcing the message just show what it says now. + // Never relayed to the game server, and no new notification: a push still + // being held for the message just says what it says now. private async announceEdit( type: ChatLobbyType, id: string, @@ -1619,11 +1620,8 @@ export class ChatService { this.logger.warn(`unable to broadcast an edit to ${type}:${id}`, error); }); - await this.notifications - .updateChatMessagePreview( - messageId, - ChatService.notificationPreview(text), - ) + await this.pushNotifications + .editChatMessage(messageId, ChatService.notificationPreview(text)) .catch((error) => { this.logger.warn( `unable to update notifications for ${type}:${id} message ${messageId}`, @@ -1639,12 +1637,14 @@ export class ChatService { id: string, messageId: string, ) { - await this.notifications.retractChatMessage(messageId).catch((error) => { - this.logger.warn( - `unable to retract notifications for ${type}:${id} message ${messageId}`, - error, - ); - }); + await this.pushNotifications + .retractChatMessage(messageId) + .catch((error) => { + this.logger.warn( + `unable to retract notifications for ${type}:${id} message ${messageId}`, + error, + ); + }); } private async recordDeletion( @@ -1765,8 +1765,9 @@ export class ChatService { // What replaced it is a signal that means what it says: the client reports // the thread it is showing while visible, and the recipient's read cursor // says how far they have got. Both are checked at send time rather than here, - // because between an insert and a push is precisely when someone opens the - // conversation. See notifications/push/push-notifications.service.ts. + // because a burst's summary goes out seconds after its messages, which is + // precisely when someone opens the conversation. See + // notifications/push/push-notifications.service.ts. private async notifyLobbyMembers( type: ChatLobbyType, id: string, @@ -1790,92 +1791,18 @@ export class ChatService { return; } - const entityId = `${type}:${id}`; - const notificationType = ChatService.notificationTypeFor(type); - - await this.notifications.notifyPlayers(notificationType, { + await this.pushNotifications.sendChatMessage(targets, { + messageId, + type: ChatService.notificationTypeFor(type), title: senderName, message: ChatService.notificationPreview(message), - role: "user", - entity_id: entityId, - steamIds: targets, - data: { - threadKey: chatThreadKey(type, id), - threadLabel: await this.threadLabel(type, id, sender), - icon: sender.avatar_url, - senderSteamId, - messageId, - }, + entityId: `${type}:${id}`, + threadKey: chatThreadKey(type, id), + threadLabel: await this.threadLabel(type, id, sender), + icon: sender.avatar_url, + senderSteamId, + blockExemptRoles: ChatService.blockExemptRoles(type), }); - - // Collapse to one unread bell row per conversation, and only after the - // insert above. Editing an existing row instead would produce no INSERT, - // and the INSERT is what the push event trigger fires on -- so the bell - // would be tidy and the phone would stay silent. - await this.notifications.collapseOlderUnread( - notificationType, - entityId, - targets, - ); - - await this.catchUpNotifications(type, id, messageId); - } - - // A delete or an edit that landed while the rows were being written had no - // rows to act on yet, so whatever became of the message is applied now. - // - // A room's audit row is what says it was deleted: unlike the redis field, it - // is not gone just because the message expired or moved. A direct message - // has no audit, but one written moments ago cannot have been pruned yet, so - // a missing row was deleted. - private async catchUpNotifications( - type: ChatLobbyType, - id: string, - messageId: string, - ) { - await this.notifications.retractChatMessageFromBlocked( - messageId, - ChatService.blockExemptRoles(type), - ); - - if (type === ChatLobbyType.Direct) { - const [row] = await this.postgres.query< - Array<{ message: string; edited_at: Date | null }> - >( - `SELECT message, edited_at FROM public.direct_messages - WHERE id = $1::uuid`, - [messageId], - ); - - if (!row) { - await this.notifications.retractChatMessage(messageId); - return; - } - - if (row.edited_at) { - await this.notifications.updateChatMessagePreview( - messageId, - ChatService.notificationPreview(row.message), - ); - } - - return; - } - - if (await this.wasDeleted(messageId)) { - await this.notifications.retractChatMessage(messageId); - return; - } - - const raw = await this.redis.hget(`chat_${type}_${id}`, messageId); - const message = raw ? (JSON.parse(raw) as ChatMessage) : null; - - if (message?.edited_at) { - await this.notifications.updateChatMessagePreview( - messageId, - ChatService.notificationPreview(message.message), - ); - } } private static notificationPreview(message: string) { @@ -1884,18 +1811,6 @@ export class ChatService { ); } - private async wasDeleted(messageId: string): Promise { - const [row] = await this.postgres.query>( - `SELECT EXISTS ( - SELECT 1 FROM public.chat_message_deletions - WHERE message_id = $1::uuid - ) AS deleted`, - [messageId], - ); - - return row?.deleted === true; - } - // Match chat is its own notification type, and so its own push category. // // Every line typed in-game is relayed -- all chat into the match room by @@ -1904,9 +1819,6 @@ export class ChatService { // line, and the player it reaches is the one already reading those lines in // the game. Sharing a category with direct messages meant the only way to // stop that was to mute DMs too. - // - // The insert, the bell collapse and the read-clear all have to agree on the - // type or the collapse stops collapsing and the badge never clears. public static notificationTypeFor( type: ChatLobbyType, ): e_notification_types_enum { @@ -2175,9 +2087,8 @@ export class ChatService { // notifyLobbyMembers bailed on the empty list -- the organizers' room // has never notified anyone in it. // - // Not narrowed to recently active staff. The list is small, and - // notifyPlayers already drops anyone with neither the bell nor a - // subscription to deliver to. + // Not narrowed to recently active staff. The list is small, and the + // push only reaches those with a device subscribed. for (const steamId of await this.organizerSteamIds()) { add(steamId); } @@ -2508,12 +2419,6 @@ export class ChatService { [user.steam_id, thread], ); - await this.notifications.markConversationRead( - ChatService.notificationTypeFor(type), - `${type}:${id}`, - user.steam_id, - ); - if (!row) { return null; } @@ -2611,9 +2516,8 @@ export class ChatService { // // Cursors rather than counts, deliberately. The client is handed a room's // whole history when it joins, so it can count what is newer than the cursor - // itself -- and counting server side would mean either reaching into every - // lobby's redis hash on page load, or reading the bell, where - // collapseOlderUnread has already reduced each conversation to one row. + // itself -- and counting server side would mean reaching into every lobby's + // redis hash on page load. public async getReadState(user: User) { const rows = await this.postgres.query< Array<{ thread: string; last_read_at: Date }> diff --git a/src/notifications/notifications.service.ts b/src/notifications/notifications.service.ts index 7211dec3..7fe4701c 100644 --- a/src/notifications/notifications.service.ts +++ b/src/notifications/notifications.service.ts @@ -760,120 +760,6 @@ export class NotificationsService { ); } - // Soft-deletes every unread row for an entity except the newest, so a busy - // conversation shows one bell entry rather than one per message. Push is - // unaffected -- each message already fired its own INSERT. - async collapseOlderUnread( - type: e_notification_types_enum, - entityId: string, - steamIds: string[], - ) { - if (steamIds.length === 0) { - return; - } - - await this.postgres.query( - `UPDATE public.notifications n - SET deleted_at = now() - WHERE n.type = $1 - AND n.entity_id = $2 - AND n.steam_id = ANY($3::bigint[]) - AND n.is_read = false - AND n.deleted_at IS NULL - AND n.id <> ( - SELECT newest.id - FROM public.notifications newest - WHERE newest.type = n.type - AND newest.entity_id = n.entity_id - AND newest.steam_id = n.steam_id - AND newest.deleted_at IS NULL - ORDER BY newest.created_at DESC - LIMIT 1 - )`, - [type, entityId, steamIds], - ); - } - - // By message id alone: a draft lobby's history moves into the match room with - // its ids intact, while its rows keep the draft's type and entity. - // - // Soft-deleting is also what stops a push still waiting in a bundling window, - // since chat delivery skips deleted rows. The text goes too, because a - // recipient can read and restore their own deleted rows. Rows - // collapseOlderUnread already retired stay retired, even when the one that - // superseded them goes. - async retractChatMessage(messageId: string) { - await this.postgres.query( - `UPDATE public.notifications - SET deleted_at = COALESCE(deleted_at, now()), - message = '' - WHERE data->>'messageId' = $1 - AND type IN ('ChatMessage', 'MatchChatMessage') - AND (deleted_at IS NULL OR message <> '')`, - [messageId], - ); - } - - // For a recipient who blocked the sender after the rows were aimed but - // before they were written, which the block's own cleanup ran too early to - // see. Blanked the way retractChatMessage blanks, for the same reasons. - async retractChatMessageFromBlocked( - messageId: string, - exemptRoles: Array = [], - ) { - await this.postgres.query( - `UPDATE public.notifications n - SET deleted_at = COALESCE(n.deleted_at, now()), - message = '' - WHERE n.data->>'messageId' = $1 - AND n.type IN ('ChatMessage', 'MatchChatMessage') - AND (n.deleted_at IS NULL OR n.message <> '') - AND EXISTS ( - SELECT 1 - FROM public.player_blocks pb - JOIN public.players p ON p.steam_id = pb.blocker_steam_id - WHERE pb.blocker_steam_id = n.steam_id - AND pb.blocked_steam_id::text = n.data->>'senderSteamId' - AND p.role::text <> ALL($2::text[]) - )`, - [messageId, exemptRoles], - ); - } - - // Read and collapsed rows too: the bell keeps read rows on show, and a - // recipient can restore a collapsed one, so either would otherwise keep the - // text the author took back. A preview is never empty, which is what tells a - // row retractChatMessage blanked apart -- an edit racing a delete must not - // write text back into it. - async updateChatMessagePreview(messageId: string, preview: string) { - await this.postgres.query( - `UPDATE public.notifications - SET message = $2 - WHERE data->>'messageId' = $1 - AND type IN ('ChatMessage', 'MatchChatMessage') - AND message <> ''`, - [messageId, preview], - ); - } - - // Opening a conversation should clear its badge everywhere, not just in the - // tab that was open. - async markConversationRead( - type: e_notification_types_enum, - entityId: string, - steamId: string, - ) { - await this.postgres.query( - `UPDATE public.notifications - SET is_read = true - WHERE type = $1 - AND entity_id = $2 - AND steam_id = $3::bigint - AND is_read = false`, - [type, entityId, steamId], - ); - } - // The poster for a map, as a push notification image. Posters live in the // web bundle (`/img/maps/screenshots/...`), not behind the API, so they are // qualified here rather than by the push service's API-relative rule. diff --git a/src/notifications/preferences/notification-categories.spec.ts b/src/notifications/preferences/notification-categories.spec.ts index 7df605f1..e2b6a1fd 100644 --- a/src/notifications/preferences/notification-categories.spec.ts +++ b/src/notifications/preferences/notification-categories.spec.ts @@ -116,6 +116,11 @@ describe("notification categories", () => { expect(offenders).toEqual([]); }); + it("offers no in-app toggle for chat, which never reaches the bell", () => { + expect(inAppKeyForType("ChatMessage")).toBeNull(); + expect(inAppKeyForType("MatchChatMessage")).toBeNull(); + }); + it("resolves in-app keys back to their own type", () => { for (const entry of IN_APP_KEYS) { expect(inAppKeyForType(entry.key)).toEqual(entry); diff --git a/src/notifications/preferences/notification-categories.ts b/src/notifications/preferences/notification-categories.ts index 0cba34cf..010e96cc 100644 --- a/src/notifications/preferences/notification-categories.ts +++ b/src/notifications/preferences/notification-categories.ts @@ -116,8 +116,6 @@ export const PUSH_KEYS: PreferenceKey[] = [ // insert time against a known recipient list, and a role-broadcast row has no // such list to filter against. export const IN_APP_KEYS: PreferenceKey[] = [ - { key: "ChatMessage", defaultEnabled: true }, - { key: "MatchChatMessage", defaultEnabled: true }, { key: "TeamInvite", defaultEnabled: true }, { key: "TournamentTeamInvite", defaultEnabled: true }, { key: "TournamentInvite", defaultEnabled: true }, diff --git a/src/notifications/push/notification-delivery.ts b/src/notifications/push/notification-delivery.ts index 33a523ac..7654d121 100644 --- a/src/notifications/push/notification-delivery.ts +++ b/src/notifications/push/notification-delivery.ts @@ -8,8 +8,8 @@ export type DeliveryPolicy = { // 0 disables windowing entirely -- one push per row, as before. bundleSeconds: number; // Drop the push if the bell row has been read or dismissed by the time it is - // sent. For chat this also means "the recipient's read cursor for the thread - // has moved past this message". + // sent. Chat has no row; for it this means "the recipient's read cursor for + // the thread has moved past this message". requireUnseen: boolean; // How long the push service may hold the message for an offline device // before discarding it. Unset keeps web-push's four-week default. @@ -183,6 +183,10 @@ export function chatThreadKey(type: string, id: string): string { return `chat:${type}:${id}`; } +export function isChatThreadKey(thread: string): boolean { + return thread.startsWith("chat:"); +} + // Where a client's focus reports land. Read by the delivery gate and written by // the socket gateway, which is why it lives here with the key it holds rather // than on either side of that exchange -- the two modules would otherwise have diff --git a/src/notifications/push/push-notifications.service.spec.ts b/src/notifications/push/push-notifications.service.spec.ts index 7c47a278..52f7e904 100644 --- a/src/notifications/push/push-notifications.service.spec.ts +++ b/src/notifications/push/push-notifications.service.spec.ts @@ -39,6 +39,31 @@ const chainableMulti = (result: unknown[]) => { return multi as any; }; +const ROOM = "chat:match:m-1"; + +const chatPush = (overrides: Record = {}) => ({ + messageId: "11111111-1111-1111-1111-111111111111", + type: "ChatMessage" as const, + title: "Luke", + message: "hey", + entityId: "match:m-1", + threadKey: ROOM, + threadLabel: "Ancients vs Ratz", + icon: null as string | null, + senderSteamId: "76561100000000009", + blockExemptRoles: [] as any[], + ...overrides, +}); + +// What a window holds for a chat message: the push itself, stamped with when +// postgres says it was queued. +const heldChat = (overrides: Record = {}) => + JSON.stringify({ + ...chatPush(), + at: "2026-10-02T12:00:00.000Z", + ...overrides, + }); + const subscription = (id: string) => ({ id, endpoint: `https://fcm.googleapis.com/fcm/send/${id}`, @@ -56,6 +81,8 @@ describe("PushNotificationsService", () => { // looking at anything, which is the case every pre-existing test assumes. let focus: Record; let pipelined: string[]; + // Deleted and edited chat messages, as the keys a delete or an edit writes. + let overrides: Record; const redis = { exists: jest.fn().mockResolvedValue(0), subscribe: jest.fn().mockResolvedValue(1), @@ -66,6 +93,9 @@ describe("PushNotificationsService", () => { get: jest.fn().mockResolvedValue(null), ttl: jest.fn().mockResolvedValue(-2), del: jest.fn().mockResolvedValue(1), + mget: jest.fn(async (...keys: string[]) => + keys.map((key) => overrides[key] ?? null), + ), rpush: jest.fn().mockResolvedValue(1), expire: jest.fn().mockResolvedValue(1), multi: jest.fn(() => chainableMulti([[null, []]])), @@ -114,6 +144,10 @@ describe("PushNotificationsService", () => { let updates: Array<{ sql: string; bindings: any[] }>; // What the badge-count query answers with. let unread: number; + // The recipient's read cursor for the thread a closing chat window is on. + let readTo: string | null; + // Senders the recipient has blocked, among those a chat window holds. + let blockedSenders: string[]; // Keys are resolved from settings (with env taking precedence), so they are // not known until loadKeys() runs. @@ -158,6 +192,9 @@ describe("PushNotificationsService", () => { updates = []; settings = {}; unread = 4; + overrides = {}; + readTo = null; + blockedSenders = []; postgres.query.mockImplementation(async (sql: string, bindings: any[]) => { if (sql.includes("AS unread")) { @@ -166,6 +203,32 @@ describe("PushNotificationsService", () => { if (sql.includes("FROM public.notifications\n")) { return notificationRow ? [notificationRow] : []; } + if (sql.includes("now() AS at")) { + return ( + recipients.length > 0 ? recipients : ["76561100000000001"] + ).flatMap((steam_id) => + subscriptions.map((sub) => ({ + steam_id, + quiet_seconds: quietSeconds, + subscription_id: sub.id, + endpoint: sub.endpoint, + p256dh: sub.p256dh, + auth: sub.auth, + at: new Date("2026-10-02T12:00:00.000Z"), + })), + ); + } + if (sql.includes("chat_read_state crs")) { + return subscriptions.map((sub) => ({ + quiet_seconds: quietSeconds, + last_read_at: readTo ? new Date(readTo) : null, + blocked: blockedSenders, + subscription_id: sub.id, + endpoint: sub.endpoint, + p256dh: sub.p256dh, + auth: sub.auth, + })); + } if (sql.includes("push_subscriptions ps")) { return bundled.length > 0 ? bundled : deliveryRows(); } @@ -413,76 +476,224 @@ describe("PushNotificationsService", () => { }); describe("focus gating", () => { - const chat = () => - notification({ - type: "ChatMessage", - title: "Luke", - message: "hey", - entity_id: "match:m-1", - data: { threadKey: "chat:match:m-1", threadLabel: "Ancients vs Ratz" }, - }); - it("says nothing to a player already reading the conversation", async () => { - notificationRow = chat(); - recipients = ["76561100000000001"]; - focus["76561100000000001"] = ["chat:match:m-1"]; + focus["76561100000000001"] = [ROOM]; - await service.sendForNotification({ - id: notificationRow.id, - type: "ChatMessage", - }); + await service.sendChatMessage(["76561100000000001"], chatPush()); expect(webPush.sendNotification).not.toHaveBeenCalled(); }); it("still buzzes a player looking at a different conversation", async () => { - notificationRow = chat(); - recipients = ["76561100000000001"]; focus["76561100000000001"] = ["chat:direct:1:2"]; - await service.sendForNotification({ - id: notificationRow.id, - type: "ChatMessage", - }); + await service.sendChatMessage(["76561100000000001"], chatPush()); expect(webPush.sendNotification).toHaveBeenCalledTimes(1); }); - it("gates each recipient of a fan-out on their own attention", async () => { - // One row per lobby member, so the reader and the absentee arrive - // together and only one of them should hear about it. - notificationRow = chat(); + it("gates each recipient of a lobby on their own attention", async () => { recipients = ["76561100000000001", "76561100000000002"]; - focus["76561100000000001"] = ["chat:match:m-1"]; + focus["76561100000000001"] = [ROOM]; - await service.sendForIds([notificationRow.id]); + await service.sendChatMessage(recipients, chatPush()); expect(webPush.sendNotification).toHaveBeenCalledTimes(1); }); }); - describe("bundling", () => { - const chat = (overrides: Record = {}) => - notification({ - type: "ChatMessage", - title: "Luke", - message: "hey", - entity_id: "match:m-1", - data: { threadKey: "chat:match:m-1", threadLabel: "Ancients vs Ratz" }, - ...overrides, + describe("chat", () => { + const payloadOf = (call: number) => + JSON.parse((webPush.sendNotification as jest.Mock).mock.calls[call][1]); + + const closeWindow = async (entries: string[]) => { + redis.multi.mockReturnValueOnce(chainableMulti([[null, entries]])); + await service.sendPending("76561100000000001", ROOM); + }; + + it("never reads or writes a notifications row for the message", async () => { + await service.sendChatMessage(["76561100000000001"], chatPush()); + await closeWindow([heldChat()]); + + expect(webPush.sendNotification).toHaveBeenCalledTimes(2); + // The badge still counts the bell; nothing else touches it. + expect( + postgres.query.mock.calls.filter( + ([sql]: [string]) => + sql.includes("public.notifications") && !sql.includes("AS unread"), + ), + ).toEqual([]); + }); + + it("asks only the push category the room belongs to", async () => { + await service.sendChatMessage( + ["76561100000000001"], + chatPush({ type: "MatchChatMessage" }), + ); + + const [, bindings] = postgres.query.mock.calls.find(([sql]: [string]) => + sql.includes("now() AS at"), + ); + + // match_chat is off unless the player turned it on. + expect(bindings).toEqual([["76561100000000001"], "match_chat", false]); + }); + + it("does not push a message deleted before it went out", async () => { + await service.retractChatMessage(chatPush().messageId); + overrides[redis.set.mock.calls.at(-1)[0]] = "1"; + + await service.sendChatMessage(["76561100000000001"], chatPush()); + + expect(webPush.sendNotification).not.toHaveBeenCalled(); + }); + + it("pushes what an edit made of the message", async () => { + await service.editChatMessage(chatPush().messageId, "fixed"); + overrides[redis.set.mock.calls.at(-1)[0]] = "fixed"; + + await service.sendChatMessage(["76561100000000001"], chatPush()); + + expect(payloadOf(0)).toMatchObject({ body: "fixed" }); + }); + + it("keeps deletes and edits for as long as a night of quiet hours", async () => { + await service.retractChatMessage("m-a"); + await service.editChatMessage("m-b", "fixed"); + + expect(redis.set).toHaveBeenCalledWith( + "notifications:chat-retracted:m-a", + 1, + "EX", + 25 * 60 * 60, + ); + expect(redis.set).toHaveBeenCalledWith( + "notifications:chat-edited:m-b", + "fixed", + "EX", + 25 * 60 * 60, + ); + }); + + it("leaves a deleted message out of the summary", async () => { + overrides["notifications:chat-retracted:m-b"] = "1"; + + await closeWindow([ + heldChat({ messageId: "m-a" }), + heldChat({ messageId: "m-b" }), + heldChat({ messageId: "m-c" }), + ]); + + expect(payloadOf(0)).toMatchObject({ body: "2 new messages" }); + }); + + it("does not repeat the first message when everything after it was deleted", async () => { + overrides["notifications:chat-retracted:m-b"] = "1"; + + await closeWindow([ + heldChat({ messageId: "m-a", pushed: true }), + heldChat({ messageId: "m-b" }), + ]); + + expect(webPush.sendNotification).not.toHaveBeenCalled(); + }); + + it("still delivers a lone message quiet hours held back", async () => { + await closeWindow([heldChat({ messageId: "m-a" })]); + + expect(payloadOf(0)).toMatchObject({ body: "hey" }); + }); + + it("says nothing when every held message was deleted", async () => { + overrides["notifications:chat-retracted:m-a"] = "1"; + + await closeWindow([heldChat({ messageId: "m-a" })]); + + expect(webPush.sendNotification).not.toHaveBeenCalled(); + }); + + it("shows a held message as it was edited", async () => { + overrides["notifications:chat-edited:m-a"] = "fixed"; + + await closeWindow([heldChat({ messageId: "m-a", message: "typo" })]); + + expect(payloadOf(0)).toMatchObject({ body: "fixed" }); + }); + + it("counts only what the player has not read since", async () => { + readTo = "2026-10-02T12:00:05.000Z"; + + await closeWindow([ + heldChat({ messageId: "m-a", at: "2026-10-02T12:00:00.000Z" }), + heldChat({ messageId: "m-b", at: "2026-10-02T12:00:10.000Z" }), + heldChat({ messageId: "m-c", at: "2026-10-02T12:00:11.000Z" }), + ]); + + expect(payloadOf(0)).toMatchObject({ body: "2 new messages" }); + }); + + it("says nothing once the player has read the whole burst", async () => { + readTo = "2026-10-02T12:00:30.000Z"; + + await closeWindow([ + heldChat({ messageId: "m-a" }), + heldChat({ messageId: "m-b" }), + ]); + + expect(webPush.sendNotification).not.toHaveBeenCalled(); + }); + + it("drops what a sender the player blocked meanwhile said", async () => { + blockedSenders = ["76561100000000009"]; + + await closeWindow([ + heldChat({ messageId: "m-a", senderSteamId: "76561100000000009" }), + heldChat({ + messageId: "m-b", + title: "Ratz", + message: "still here", + senderSteamId: "76561100000000008", + }), + ]); + + expect(payloadOf(0)).toMatchObject({ + title: "Ratz · Ancients vs Ratz", + body: "still here", }); + }); + + it("lets a moderator's block leave a group room alone", async () => { + await closeWindow([ + heldChat({ blockExemptRoles: ["moderator", "administrator"] }), + ]); + const [, bindings] = postgres.query.mock.calls.find(([sql]: [string]) => + sql.includes("chat_read_state crs"), + ); + + expect(bindings).toEqual([ + "76561100000000001", + ROOM, + ["76561100000000009"], + ["moderator", "administrator"], + "chat", + true, + ]); + }); + + it("skips an entry it cannot read rather than the whole window", async () => { + await closeWindow(["not json", heldChat()]); + + expect(payloadOf(0)).toMatchObject({ body: "hey" }); + }); + }); + + describe("bundling", () => { const payloadOf = (call: number) => JSON.parse((webPush.sendNotification as jest.Mock).mock.calls[call][1]); it("pushes the first message of a burst straight away", async () => { - notificationRow = chat(); - recipients = ["76561100000000001"]; - - await service.sendForNotification({ - id: notificationRow.id, - type: "ChatMessage", - }); + await service.sendChatMessage(["76561100000000001"], chatPush()); expect(webPush.sendNotification).toHaveBeenCalledTimes(1); // Named room and all: a message on its own still has to say where it @@ -490,46 +701,44 @@ describe("PushNotificationsService", () => { expect(payloadOf(0)).toMatchObject({ title: "Luke · Ancients vs Ratz", body: "hey", - tag: "chat:match:m-1", + url: "/chat/match%3Am-1", + tag: ROOM, renotify: true, - threadKey: "chat:match:m-1", + threadKey: ROOM, }); }); it("does not repeat a direct message's sender as its room", async () => { // A DM's label is whoever sent it, so naming the room would say the // same name twice. - notificationRow = chat({ - entity_id: "direct:1:2", - data: { threadKey: "chat:direct:1:2", threadLabel: "Luke" }, - }); - recipients = ["76561100000000001"]; - - await service.sendForNotification({ - id: notificationRow.id, - type: "ChatMessage", - }); + await service.sendChatMessage( + ["76561100000000001"], + chatPush({ + entityId: "direct:1:2", + threadKey: "chat:direct:1:2", + threadLabel: "Luke", + }), + ); expect(payloadOf(0)).toMatchObject({ title: "Luke" }); }); it("holds a message that lands inside an open window", async () => { - notificationRow = chat(); - recipients = ["76561100000000001"]; // The window is already taken, which is what a second message sees. redis.set.mockResolvedValueOnce(null); redis.get.mockResolvedValueOnce("token-1"); redis.ttl.mockResolvedValueOnce(12); - await service.sendForNotification({ - id: notificationRow.id, - type: "ChatMessage", - }); + await service.sendChatMessage(["76561100000000001"], chatPush()); expect(webPush.sendNotification).not.toHaveBeenCalled(); + expect(redis.rpush).toHaveBeenCalledWith( + `notifications:push-pending:76561100000000001:${ROOM}`, + heldChat(), + ); expect(pushDeliveryQueue.add).toHaveBeenCalledWith( SendPushDelivery.name, - { steamId: "76561100000000001", thread: "chat:match:m-1" }, + { steamId: "76561100000000001", thread: ROOM }, expect.objectContaining({ // Colons encoded: BullMQ rejects a custom job id containing one. jobId: "push-trail.76561100000000001.chat%3Amatch%3Am-1.token-1", @@ -544,36 +753,25 @@ describe("PushNotificationsService", () => { it("counts the message that opened the window", async () => { // The summary replaces the leading notification on the device, so // leaving it out of the tally makes a burst of four report three. - notificationRow = chat(); - recipients = ["76561100000000001"]; - - await service.sendForNotification({ - id: notificationRow.id, - type: "ChatMessage", - }); + await service.sendChatMessage(["76561100000000001"], chatPush()); const reset = redis.multi.mock.results.at(-1)?.value; expect(reset.del).toHaveBeenCalledWith( - "notifications:push-pending:76561100000000001:chat:match:m-1", + `notifications:push-pending:76561100000000001:${ROOM}`, ); expect(reset.rpush).toHaveBeenCalledWith( - "notifications:push-pending:76561100000000001:chat:match:m-1", - notificationRow.id, + `notifications:push-pending:76561100000000001:${ROOM}`, + heldChat({ pushed: true }), ); }); it("takes the leading edge when the window expired mid-decision", async () => { - notificationRow = chat(); - recipients = ["76561100000000001"]; redis.set.mockResolvedValueOnce(null); redis.get.mockResolvedValueOnce(null); redis.ttl.mockResolvedValueOnce(-2); - await service.sendForNotification({ - id: notificationRow.id, - type: "ChatMessage", - }); + await service.sendChatMessage(["76561100000000001"], chatPush()); expect(webPush.sendNotification).toHaveBeenCalledTimes(1); expect(pushDeliveryQueue.add).not.toHaveBeenCalled(); @@ -583,21 +781,16 @@ describe("PushNotificationsService", () => { // Both attempts lost the race and both then found the key already gone. // Reporting a leading edge without holding the key opens a bundle that // nothing will ever drain, and the next message takes the edge as well. - notificationRow = chat(); - recipients = ["76561100000000001"]; redis.set.mockResolvedValueOnce(null).mockResolvedValueOnce(null); redis.get.mockResolvedValue(null); redis.ttl.mockResolvedValue(-2); - await service.sendForNotification({ - id: notificationRow.id, - type: "ChatMessage", - }); + await service.sendChatMessage(["76561100000000001"], chatPush()); expect(webPush.sendNotification).toHaveBeenCalledTimes(1); // Plain SET, not SET NX: the two NX attempts above are what just failed. expect(redis.set).toHaveBeenLastCalledWith( - "notifications:push-window:76561100000000001:chat:match:m-1", + `notifications:push-window:76561100000000001:${ROOM}`, expect.any(String), "EX", expect.any(Number), @@ -608,71 +801,61 @@ describe("PushNotificationsService", () => { // A bundling window opened seconds before 22:00 closes inside quiet // hours, and delivering its summary there is the buzz the hold exists to // prevent. - const held = ["id-a", "id-b"]; - redis.multi.mockReturnValueOnce(chainableMulti([[null, held]])); - - notificationRow = chat(); - bundled = held.map((id) => ({ - ...chat({ id }), - steam_id: "76561100000000001", - quiet_seconds: 6 * 60 * 60, - subscription_id: "sub-1", - endpoint: subscription("sub-1").endpoint, - p256dh: "p256dh", - auth: "auth", - })); + redis.multi.mockReturnValueOnce( + chainableMulti([ + [ + null, + [heldChat({ messageId: "m-a" }), heldChat({ messageId: "m-b" })], + ], + ]), + ); + quietSeconds = 6 * 60 * 60; - await service.sendPending("76561100000000001", "chat:match:m-1"); + await service.sendPending("76561100000000001", ROOM); expect(webPush.sendNotification).not.toHaveBeenCalled(); expect(pushDeliveryQueue.add).toHaveBeenCalledWith( SendPushDelivery.name, - { steamId: "76561100000000001", thread: "chat:match:m-1" }, + { steamId: "76561100000000001", thread: ROOM }, expect.objectContaining({ delay: 6 * 60 * 60 * 1000 }), ); }); it("replaces the burst with one summary when the window closes", async () => { - const held = ["id-a", "id-b", "id-c"]; - redis.multi.mockReturnValueOnce(chainableMulti([[null, held]])); - - notificationRow = chat(); - bundled = held.map((id) => ({ - ...chat({ id }), - steam_id: "76561100000000001", - subscription_id: "sub-1", - endpoint: subscription("sub-1").endpoint, - p256dh: "p256dh", - auth: "auth", - })); + redis.multi.mockReturnValueOnce( + chainableMulti([ + [ + null, + ["m-a", "m-b", "m-c"].map((messageId) => heldChat({ messageId })), + ], + ]), + ); - await service.sendPending("76561100000000001", "chat:match:m-1"); + await service.sendPending("76561100000000001", ROOM); expect(webPush.sendNotification).toHaveBeenCalledTimes(1); expect(payloadOf(0)).toMatchObject({ title: "Luke · Ancients vs Ratz", body: "3 new messages", - tag: "chat:match:m-1", + tag: ROOM, renotify: true, - threadKey: "chat:match:m-1", + threadKey: ROOM, }); }); it("names the room when a burst has more than one sender", async () => { - const held = ["id-a", "id-b", "id-c"]; - redis.multi.mockReturnValueOnce(chainableMulti([[null, held]])); - - notificationRow = chat(); - bundled = ["Luke", "Ratz", "Catz"].map((title, index) => ({ - ...chat({ id: held[index], title }), - steam_id: "76561100000000001", - subscription_id: "sub-1", - endpoint: subscription("sub-1").endpoint, - p256dh: "p256dh", - auth: "auth", - })); + redis.multi.mockReturnValueOnce( + chainableMulti([ + [ + null, + ["Luke", "Ratz", "Catz"].map((title, index) => + heldChat({ messageId: `m-${index}`, title }), + ), + ], + ]), + ); - await service.sendPending("76561100000000001", "chat:match:m-1"); + await service.sendPending("76561100000000001", ROOM); expect(payloadOf(0)).toMatchObject({ title: "Ancients vs Ratz", @@ -683,19 +866,14 @@ describe("PushNotificationsService", () => { it("holds a message until quiet hours are over", async () => { // Dropped outright before, so a night of messages arrived as nothing at // all -- a silent phone and a full bell in the morning. - notificationRow = chat(); - recipients = ["76561100000000001"]; quietSeconds = 6 * 60 * 60; - await service.sendForNotification({ - id: notificationRow.id, - type: "ChatMessage", - }); + await service.sendChatMessage(["76561100000000001"], chatPush()); expect(webPush.sendNotification).not.toHaveBeenCalled(); expect(pushDeliveryQueue.add).toHaveBeenCalledWith( SendPushDelivery.name, - { steamId: "76561100000000001", thread: "chat:match:m-1" }, + { steamId: "76561100000000001", thread: ROOM }, // Woken when the window closes, not on the bundling window. expect.objectContaining({ delay: 6 * 60 * 60 * 1000 }), ); @@ -718,19 +896,14 @@ describe("PushNotificationsService", () => { it("keeps the held payload alive past the whole quiet window", async () => { // The pending list used to expire after fifteen minutes, which would // have thrown the night away long before anyone woke up. - notificationRow = chat(); - recipients = ["76561100000000001"]; quietSeconds = 8 * 60 * 60; - await service.sendForNotification({ - id: notificationRow.id, - type: "ChatMessage", - }); + await service.sendChatMessage(["76561100000000001"], chatPush()); const reset = redis.multi.mock.results.at(-1)?.value; expect(reset.expire).toHaveBeenCalledWith( - "notifications:push-pending:76561100000000001:chat:match:m-1", + `notifications:push-pending:76561100000000001:${ROOM}`, 8 * 60 * 60 + 300, ); }); @@ -738,32 +911,51 @@ describe("PushNotificationsService", () => { it("releases the window even when nothing survives the gate", async () => { // Otherwise the next burst's leading push is swallowed too, and goes on // being swallowed until the key expires on its own. - await service.sendPending("76561100000000001", "chat:match:m-1"); + await service.sendPending("76561100000000001", ROOM); expect(redis.del).toHaveBeenCalledWith( - "notifications:push-window:76561100000000001:chat:match:m-1", + `notifications:push-window:76561100000000001:${ROOM}`, ); expect(webPush.sendNotification).not.toHaveBeenCalled(); }); it("says nothing if the player opened the thread while it was held", async () => { + redis.multi.mockReturnValueOnce( + chainableMulti([ + [ + null, + [heldChat({ messageId: "m-a" }), heldChat({ messageId: "m-b" })], + ], + ]), + ); + focus["76561100000000001"] = [ROOM]; + + await service.sendPending("76561100000000001", ROOM); + + expect(webPush.sendNotification).not.toHaveBeenCalled(); + }); + + it("still summarises a bundled type that writes rows", async () => { const held = ["id-a", "id-b"]; redis.multi.mockReturnValueOnce(chainableMulti([[null, held]])); - notificationRow = chat(); + notificationRow = notification(); bundled = held.map((id) => ({ - ...chat({ id }), + ...notification({ id }), steam_id: "76561100000000001", + quiet_seconds: 0, subscription_id: "sub-1", endpoint: subscription("sub-1").endpoint, p256dh: "p256dh", auth: "auth", })); - focus["76561100000000001"] = ["chat:match:m-1"]; - await service.sendPending("76561100000000001", "chat:match:m-1"); + await service.sendPending("76561100000000001", "MatchStatusChange:m-1"); - expect(webPush.sendNotification).not.toHaveBeenCalled(); + expect(payloadOf(0)).toMatchObject({ + title: "Match ready", + body: "2 new notifications", + }); }); }); @@ -969,18 +1161,11 @@ describe("PushNotificationsService", () => { }); it("gives chat no buttons", async () => { - // Reading the bell row would leave the conversation's own cursor where - // it was, so a Dismiss here would lie. - notificationRow = notification({ - type: "ChatMessage", - title: "Luke", - message: "hey", - entity_id: "match:m-1", - data: { threadKey: "chat:match:m-1", threadLabel: "Ancients vs Ratz" }, - }); - recipients = ["76561100000000001"]; + // Chat has no bell row to mark, and marking one read would leave the + // conversation's own cursor where it was. + await service.sendChatMessage(["76561100000000001"], chatPush()); - expect((await send()).actions).toEqual([]); + expect(payloadOf(0).actions).toEqual([]); }); it("keeps a broken stored action from blocking the push", async () => { diff --git a/src/notifications/push/push-notifications.service.ts b/src/notifications/push/push-notifications.service.ts index 3be5e9bc..d86e0c1d 100644 --- a/src/notifications/push/push-notifications.service.ts +++ b/src/notifications/push/push-notifications.service.ts @@ -9,7 +9,10 @@ import { PostgresService } from "../../postgres/postgres.service"; import { RedisManagerService } from "../../redis/redis-manager/redis-manager.service"; import { AppConfig } from "src/configs/types/AppConfig"; import { WebPushConfig } from "src/configs/types/WebPushConfig"; -import { e_player_roles_enum } from "generated/schema"; +import { + e_notification_types_enum, + e_player_roles_enum, +} from "generated/schema"; import { generateMutationOp } from "../../../generated"; import { rolesAtOrAbove } from "src/utilities/isRoleAbove"; import { SystemSettingName } from "src/system/enums/SystemSettingName"; @@ -21,6 +24,7 @@ import { DEFAULT_DELIVERY_POLICY, DeliveryPolicy, deliveryPolicyForType, + isChatThreadKey, presenceFocusKey, threadKeyFor, } from "./notification-delivery"; @@ -66,6 +70,29 @@ export type NotificationRow = { actions?: NotificationAction[] | null; }; +// A chat message on its way to a phone. Chat writes no notifications row, so +// this is everything the gate knows about it, and what a bundling window holds +// until it closes. +export type ChatPush = { + messageId: string; + type: e_notification_types_enum; + title: string; + message: string; + entityId: string; + threadKey: string; + threadLabel: string; + icon?: string | null; + senderSteamId: string; + // Recipients whose block on the sender does not hide the message from them + // (ChatService.blockExemptRoles). + blockExemptRoles: e_player_roles_enum[]; +}; + +// `at` is when it was queued by postgres's clock, the one the read cursor is +// stamped with. `pushed` marks the message that opened the window, which the +// device already shows. +type HeldChatPush = ChatPush & { at: string; pushed?: boolean }; + // A notification button as the service worker sees it: an id to match the // click against and a ready-to-POST GraphQL operation, so the worker never // has to know how to build one. @@ -103,6 +130,9 @@ type Delivery = { // one. Held rather than dropped, so a night of messages arrives as one // summary in the morning instead of as nothing at all. quietSeconds: number; + // What the window holds for a summary when it is not the rows' ids, before + // and after the leading push has gone out. + held?: { waiting: string[]; sent: string[] }; }; // Which rows to consider. `ids` is the exact set a writer just inserted; @@ -205,6 +235,15 @@ const pendingKey = (steamId: string, thread: string) => const pendingTtlFor = (windowSeconds: number) => Math.max(900, windowSeconds + 300); +const chatRetractedKey = (messageId: string) => + `notifications:chat-retracted:${messageId}`; + +const chatEditedKey = (messageId: string) => + `notifications:chat-edited:${messageId}`; + +// Outlives the longest a chat push is ever held: a whole night of quiet hours. +const CHAT_OVERRIDE_TTL_SECONDS = 25 * 60 * 60; + @Injectable() export class PushNotificationsService { // How far back a batched send resolves the rows of one burst, and so how wide @@ -602,6 +641,11 @@ export class PushNotificationsService { return; } + if (isChatThreadKey(thread)) { + await this.sendPendingChat(steamId, thread, ids); + return; + } + const newest = await this.newestOf({ ids }); if (!newest) { @@ -650,14 +694,9 @@ export class PushNotificationsService { continue; } - // Counted from the window rather than from the rows that survived it. - // - // collapseOlderUnread soft-deletes every superseded ChatMessage row so - // the bell shows one entry per conversation, and requireUnseen drops - // soft-deleted rows -- so resolving a burst of four finds one survivor. - // The window is what actually knows how many arrived; the surviving rows - // are only there to say whether it is still worth sending at all, and to - // supply the text. + // Counted from the window rather than from the rows that survived it: the + // surviving rows are only there to say whether it is still worth sending + // at all, and to supply the text. await this.deliver( delivery.steamId, delivery.subscriptions, @@ -667,6 +706,260 @@ export class PushNotificationsService { } } + // A chat message, pushed straight from the conversation: there is no row + // behind it, and every check a row would have made is made here instead. + public async sendChatMessage( + steamIds: string[], + push: ChatPush, + ): Promise { + if (!this.configured || steamIds.length === 0) { + return; + } + + // Deleted or edited in the moments it took to get here. + const [current] = await this.applyChatOverrides([push]); + + if (!current) { + return; + } + + const category = pushCategoryForType(current.type); + const rows = await this.postgres.query< + Array & { at: Date }> + >( + `SELECT p.steam_id::text AS steam_id, + public.quiet_hours_seconds_remaining( + p.quiet_hours_start, p.quiet_hours_end, p.notification_timezone + ) AS quiet_seconds, + ps.id::text AS subscription_id, ps.endpoint, ps.p256dh, ps.auth, + now() AS at + FROM public.players p + JOIN public.push_subscriptions ps ON ps.steam_id = p.steam_id + LEFT JOIN public.notification_preferences np + ON np.steam_id = p.steam_id + AND np.channel = 'push' + AND np.key = $2 + WHERE p.steam_id = ANY($1::bigint[]) + AND COALESCE(np.enabled, $3::boolean) = true`, + [steamIds, category?.key ?? "", category?.defaultEnabled ?? true], + ); + + if (rows.length === 0) { + return; + } + + const row = PushNotificationsService.chatRow(current); + const held: HeldChatPush = { + ...current, + at: new Date(rows[0].at).toISOString(), + }; + const deliveries = PushNotificationsService.groupDeliveries( + rows.map((recipient) => ({ ...row, ...recipient })), + ).map((delivery) => ({ + ...delivery, + held: { + waiting: [JSON.stringify(held)], + sent: [JSON.stringify({ ...held, pushed: true })], + }, + })); + + await this.dispatch( + deliveries, + deliveryPolicyForType(current.type) ?? DEFAULT_DELIVERY_POLICY, + ); + } + + // Keyed by message id alone, because that is all a delete knows, and read by + // whichever window is holding the message whenever it closes. + public async retractChatMessage(messageId: string): Promise { + await this.redis.set( + chatRetractedKey(messageId), + 1, + "EX", + CHAT_OVERRIDE_TTL_SECONDS, + ); + } + + public async editChatMessage( + messageId: string, + preview: string, + ): Promise { + await this.redis.set( + chatEditedKey(messageId), + preview, + "EX", + CHAT_OVERRIDE_TTL_SECONDS, + ); + } + + // A chat window closing. Everything it held is checked again against what + // happened in the meantime: deleted, edited, read, or its sender blocked. + private async sendPendingChat( + steamId: string, + thread: string, + entries: string[], + ): Promise { + const held = await this.applyChatOverrides( + entries.flatMap((entry): HeldChatPush[] => { + try { + return [JSON.parse(entry) as HeldChatPush]; + } catch { + return []; + } + }), + ); + + if (held.length === 0) { + return; + } + + const newest = held.at(-1); + const policy = + deliveryPolicyForType(newest.type) ?? DEFAULT_DELIVERY_POLICY; + const category = pushCategoryForType(newest.type); + + const rows = await this.postgres.query< + Array< + Omit & { + subscription_id: string; + quiet_seconds: number; + last_read_at: Date | null; + blocked: string[] | null; + } + > + >( + `SELECT public.quiet_hours_seconds_remaining( + p.quiet_hours_start, p.quiet_hours_end, p.notification_timezone + ) AS quiet_seconds, + crs.last_read_at, + ARRAY( + SELECT pb.blocked_steam_id::text + FROM public.player_blocks pb + WHERE pb.blocker_steam_id = p.steam_id + AND pb.blocked_steam_id = ANY($3::bigint[]) + AND p.role::text <> ALL($4::text[]) + ) AS blocked, + ps.id::text AS subscription_id, ps.endpoint, ps.p256dh, ps.auth + FROM public.players p + JOIN public.push_subscriptions ps ON ps.steam_id = p.steam_id + LEFT JOIN public.notification_preferences np + ON np.steam_id = p.steam_id + AND np.channel = 'push' + AND np.key = $5 + LEFT JOIN public.chat_read_state crs + ON crs.steam_id = p.steam_id + AND crs.thread = $2 + WHERE p.steam_id = $1::bigint + AND COALESCE(np.enabled, $6::boolean) = true`, + [ + steamId, + thread, + [...new Set(held.map(({ senderSteamId }) => senderSteamId))], + newest.blockExemptRoles, + category?.key ?? "", + category?.defaultEnabled ?? true, + ], + ); + + if (rows.length === 0) { + return; + } + + const { quiet_seconds, last_read_at, blocked } = rows[0]; + const readTo = last_read_at ? new Date(last_read_at).getTime() : null; + const hidden = new Set(blocked ?? []); + + const unseen = held.filter( + (push) => + !hidden.has(push.senderSteamId) && + (readTo === null || new Date(push.at).getTime() > readTo), + ); + + // Only what the device already shows: everything after it was deleted or + // read, and saying it again would buzz for nothing new. + if (unseen.every(({ pushed }) => pushed)) { + return; + } + + if ((await this.filterFocusedOn([steamId], thread)).has(steamId)) { + return; + } + + const entriesLeft = unseen.map((push) => JSON.stringify(push)); + + if (Number(quiet_seconds ?? 0) > 0 && !policy.ignoreQuietHours) { + const claim = await this.claimWindow( + steamId, + thread, + Number(quiet_seconds), + ); + + if (claim.leading) { + await this.resetPending(steamId, thread, entriesLeft, claim.ttl); + } else { + await this.appendPending(steamId, thread, entriesLeft, claim.ttl); + } + + await this.scheduleTrailing(steamId, thread, claim.token, claim.ttl); + return; + } + + await this.deliver( + steamId, + rows.map(({ subscription_id, endpoint, p256dh, auth }) => ({ + id: subscription_id, + endpoint, + p256dh, + auth, + })), + unseen.map((push) => PushNotificationsService.chatRow(push)), + ); + } + + // Drops what was deleted after it was queued, and shows the rest as edited. + private async applyChatOverrides( + pushes: T[], + ): Promise { + if (pushes.length === 0) { + return []; + } + + const overrides = await this.redis.mget( + ...pushes.flatMap(({ messageId }) => [ + chatRetractedKey(messageId), + chatEditedKey(messageId), + ]), + ); + + return pushes.flatMap((push, index) => { + if (overrides[index * 2] !== null) { + return []; + } + + const edited = overrides[index * 2 + 1]; + + return [edited === null ? push : { ...push, message: edited }]; + }); + } + + private static chatRow(push: ChatPush): NotificationRow { + return { + id: push.messageId, + type: push.type, + role: "user", + title: push.title, + message: push.message, + entity_id: push.entityId, + data: { + threadKey: push.threadKey, + threadLabel: push.threadLabel, + icon: push.icon, + senderSteamId: push.senderSteamId, + messageId: push.messageId, + }, + }; + } + private static readonly SELECT_NOTIFICATION = `SELECT id::text AS id, type::text AS type, role::text AS role, title, message, entity_id, data, actions FROM public.notifications`; @@ -762,18 +1055,11 @@ export class PushNotificationsService { ON np.steam_id = p.steam_id AND np.channel = 'push' AND np.key = $${next + 2} - -- Only ever matches a row whose writer declared a thread, which today is - -- chat and nothing else. Anything without one joins to NULL and passes. - LEFT JOIN public.chat_read_state crs - ON crs.steam_id = p.steam_id - AND crs.thread = n.data->>'threadKey' WHERE ${selectorSql} AND COALESCE(np.enabled, $${next + 3}::boolean) = true -- Dealt with in the bell between the insert and now. AND ($${next + 4}::boolean = false OR (n.is_read = false AND n.deleted_at IS NULL)) - -- Or read in the conversation itself, which never touches the bell. - AND (crs.last_read_at IS NULL OR n.created_at > crs.last_read_at) ORDER BY n.created_at ASC`, [ ...selectorParams, @@ -784,6 +1070,10 @@ export class PushNotificationsService { ], ); + return PushNotificationsService.groupDeliveries(rows); + } + + private static groupDeliveries(rows: DeliveryRow[]): Delivery[] { const byRecipient = new Map(); for (const row of rows) { @@ -852,7 +1142,8 @@ export class PushNotificationsService { continue; } - const ids = delivery.notifications.map(({ id }) => id); + const ids = + delivery.held?.waiting ?? delivery.notifications.map(({ id }) => id); // Asleep. Hold everything until the window closes and let the trailing // job deliver it as one summary -- which is the same machinery bundling @@ -905,7 +1196,12 @@ export class PushNotificationsService { // device, so leaving it out would make a burst of four report three. // The list is reset rather than appended to, so a window that closed // without ever being drained cannot leak into the next one's count. - await this.resetPending(delivery.steamId, thread, ids, claim.ttl); + await this.resetPending( + delivery.steamId, + thread, + delivery.held?.sent ?? ids, + claim.ttl, + ); continue; } @@ -1182,8 +1478,8 @@ export class PushNotificationsService { const newest = notifications.at(-1); const byId = { id: { _in: notifications.map(({ id }) => id) } }; - // Marking a bell row read says nothing about the conversation's own read - // cursor, and a "Dismiss" that leaves the thread unread would mislead. + // Chat has no bell row to mark, and a "Dismiss" that leaves the thread + // unread would mislead. if (newest.type.endsWith("ChatMessage")) { return []; } diff --git a/test/chat-blocks.spec.ts b/test/chat-blocks.spec.ts index 2b62cfdc..6cac31d5 100644 --- a/test/chat-blocks.spec.ts +++ b/test/chat-blocks.spec.ts @@ -9,12 +9,10 @@ import { ChatErrorCode } from "./../src/chat/enums/ChatErrorCode"; import { ChatLobbyType } from "./../src/chat/enums/ChatLobbyTypes"; import { directRoomId } from "./../src/chat/utilities/directRoomId"; import { PlayerBlocksService } from "./../src/player-blocks/player-blocks.service"; -import { NotificationsService } from "./../src/notifications/notifications.service"; -import { NotificationPreferencesService } from "./../src/notifications/preferences/notification-preferences.service"; // A block hides what the blocked player says from the blocker in group rooms, -// and closes a DM in both directions. Against real Postgres (the block, its -// trigger, the bell) and real redis (rooms, history, fan-out). +// and closes a DM in both directions. Against real Postgres (the block and its +// trigger) and real redis (rooms, history, fan-out). describe("chat blocks (SQL-driven)", () => { let db: SqlTestDb; let postgres: PostgresService; @@ -26,7 +24,7 @@ describe("chat blocks (SQL-driven)", () => { let roster: string[]; let friendshipOverride: boolean; - let beforeBellInsert: (() => Promise) | undefined; + let pushes: Array<{ steamIds: string[]; messageId: string }>; let to: jest.SpyInstance; let notify: jest.SpyInstance; @@ -91,41 +89,6 @@ describe("chat blocks (SQL-driven)", () => { return {}; }), - mutation: jest.fn(async (mutation: any) => { - const insert = mutation?.insert_notifications; - - if (!insert) { - return {}; - } - - const hook = beforeBellInsert; - beforeBellInsert = undefined; - await hook?.(); - - const returning: Array<{ id: string }> = []; - - for (const object of insert.__args.objects) { - const [row] = await postgres.query>( - `INSERT INTO notifications - (type, title, message, role, steam_id, entity_id, in_app, data) - VALUES ($1, $2, $3, $4, $5::bigint, $6, $7, $8::jsonb) - RETURNING id::text AS id`, - [ - object.type, - object.title, - object.message, - object.role, - object.steam_id, - object.entity_id ?? null, - object.in_app ?? true, - object.data ? JSON.stringify(object.data) : null, - ], - ); - returning.push(row); - } - - return { insert_notifications: { returning } }; - }), }; beforeAll(async () => { @@ -141,20 +104,6 @@ describe("chat blocks (SQL-driven)", () => { postgres = db.postgres; fx = new Fixtures(postgres, 76561192820000000n); - const notifications = new NotificationsService( - hasura as any, - postgres, - logger as any, - { get: () => ({ webDomain: "https://example.com" }) } as any, - new NotificationPreferencesService(postgres), - { - filterSubscribed: async (): Promise => [], - claimFanOut: jest.fn(), - } as any, - { add: jest.fn() } as any, - { add: jest.fn() } as any, - ); - blocks = new PlayerBlocksService(postgres); chat = new ChatService( @@ -163,7 +112,16 @@ describe("chat blocks (SQL-driven)", () => { hasura as any, postgres, { getConnection: () => redis } as any, - notifications, + { + sendChatMessage: async ( + steamIds: string[], + push: { messageId: string }, + ) => { + pushes.push({ steamIds, messageId: push.messageId }); + }, + retractChatMessage: async () => {}, + editChatMessage: async () => {}, + } as any, blocks, ); }, 600_000); @@ -178,7 +136,6 @@ describe("chat blocks (SQL-driven)", () => { jest.restoreAllMocks(); jest.clearAllMocks(); await redis.flushall(); - await postgres.query("DELETE FROM notifications"); await postgres.query("DELETE FROM chat_message_edits"); await postgres.query("DELETE FROM chat_message_deletions"); await postgres.query("DELETE FROM direct_messages"); @@ -187,7 +144,7 @@ describe("chat blocks (SQL-driven)", () => { roster = []; friendshipOverride = false; - beforeBellInsert = undefined; + pushes = []; to = jest.spyOn(chat as any, "to"); notify = jest.spyOn(chat as any, "notifyLobbyMembers"); @@ -288,25 +245,11 @@ describe("chat blocks (SQL-driven)", () => { .map(({ steamId }) => steamId) .sort(); - const bell = (messageId: string) => - postgres.query< - Array<{ steam_id: string; message: string; deleted: boolean }> - >( - `SELECT steam_id::text AS steam_id, message, - deleted_at IS NOT NULL AS deleted - FROM notifications - WHERE data->>'messageId' = $1 - ORDER BY steam_id`, - [messageId], - ); - - const previews = async (messageId: string) => - Object.fromEntries( - (await bell(messageId)).map(({ steam_id, message }) => [ - steam_id, - message, - ]), - ); + const pushedTo = (messageId: string) => + pushes + .filter((push) => push.messageId === messageId) + .flatMap(({ steamIds }) => steamIds) + .sort(); describe("a group room", () => { let blocker: string; @@ -437,20 +380,12 @@ describe("chat blocks (SQL-driven)", () => { await block(blocker, blocked); - expect(await previews(before)).toEqual({ - [blocker]: "before", - [bystander]: "before", - }); - const after = messageIdOf(await inMatch(blocked, "after")); expect(recipientsOf(`lobby:match:${matchId}:chat`, after)).toEqual( [...roster].sort(), ); - expect(await previews(after)).toEqual({ - [blocker]: "after", - [bystander]: "after", - }); + expect(pushedTo(after)).toEqual([blocker, bystander].sort()); await chat.editMessage( ChatLobbyType.Match, @@ -500,70 +435,24 @@ describe("chat blocks (SQL-driven)", () => { ); }); - it("writes the blocker no bell row for a hidden line, and blanks the ones from before", async () => { + it("pushes the blocker nothing for a hidden line", async () => { const before = messageIdOf(await inMatch(blocked, "before")); - expect(await previews(before)).toEqual({ - [blocker]: "before", - [bystander]: "before", - }); + expect(pushedTo(before)).toEqual([blocker, bystander].sort()); await block(blocker, blocked); - expect(await bell(before)).toContainEqual({ - steam_id: blocker, - message: "", - deleted: true, - }); - expect(await bell(before)).toContainEqual({ - steam_id: bystander, - message: "before", - deleted: false, - }); - const after = messageIdOf(await inMatch(blocked, "after")); - expect(await previews(after)).toEqual({ [bystander]: "after" }); - - await chat.editMessage( - ChatLobbyType.Match, - matchId, - before, - player(blocked), - "before, edited", - ); - await settle(); - - expect(await previews(before)).toEqual({ - [blocker]: "", - [bystander]: "before, edited", - }); - }); - - it("blanks the blocker's row for a line whose rows were aimed before the block landed", async () => { - beforeBellInsert = () => block(blocker, blocked); - - const id = messageIdOf(await inMatch(blocked, "racing")); - - expect(await previews(id)).toEqual({ - [blocker]: "", - [bystander]: "racing", - }); - expect(await bell(id)).toContainEqual({ - steam_id: blocker, - message: "", - deleted: true, - }); + expect(pushedTo(after)).toEqual([bystander]); }); - it("leaves the bell alone for everyone when the blocker is the one talking", async () => { + it("still pushes everyone else when the blocker is the one talking", async () => { await block(blocker, blocked); const id = messageIdOf(await inMatch(blocker, "hello")); - expect((await bell(id)).map(({ steam_id }) => steam_id)).toEqual( - [blocked, bystander].sort(), - ); + expect(pushedTo(id)).toEqual([blocked, bystander].sort()); }); it("gives back lines that are still live once unblocked", async () => { @@ -772,7 +661,7 @@ describe("chat blocks (SQL-driven)", () => { await refusesEverything(); }); - it("refuses both sides and blanks the bell when the blocker is a moderator", async () => { + it("refuses both sides when the blocker is a moderator", async () => { await postgres.query("DELETE FROM chat_read_state"); await postgres.query( "UPDATE players SET role = 'moderator' WHERE steam_id = $1::bigint", @@ -781,24 +670,6 @@ describe("chat blocks (SQL-driven)", () => { await block(blocker, blocked); await refusesEverything(); - - expect(await bell(fromBlocked)).toEqual([ - { steam_id: blocker, message: "", deleted: true }, - ]); - }); - - it("blanks a moderator's row for a DM whose rows were aimed before their block landed", async () => { - await postgres.query( - "UPDATE players SET role = 'moderator' WHERE steam_id = $1::bigint", - [blocker], - ); - beforeBellInsert = () => block(blocker, blocked); - - const racing = await say(ChatLobbyType.Direct, room, blocked, "racing"); - - expect(await bell(racing.accepted ? racing.messageId : "")).toEqual([ - { steam_id: blocker, message: "", deleted: true }, - ]); }); it("refuses both sides on the block alone, while a friendship still reads as accepted", async () => { @@ -849,16 +720,5 @@ describe("chat blocks (SQL-driven)", () => { expect(await rail(blocker)).toEqual([room]); }); - - it("blanks the blocker's bell rows for the blocked player's messages, and never the other way", async () => { - await block(blocker, blocked); - - expect(await bell(fromBlocked)).toEqual([ - { steam_id: blocker, message: "", deleted: true }, - ]); - expect(await bell(fromBlocker)).toEqual([ - { steam_id: blocked, message: "hi", deleted: false }, - ]); - }); }); }); diff --git a/test/chat-direct-messages.spec.ts b/test/chat-direct-messages.spec.ts index 35cd6db1..318a9a21 100644 --- a/test/chat-direct-messages.spec.ts +++ b/test/chat-direct-messages.spec.ts @@ -7,7 +7,6 @@ import { ChatService } from "./../src/chat/chat.service"; import { PlayerBlocksService } from "./../src/player-blocks/player-blocks.service"; import { ChatErrorCode } from "./../src/chat/enums/ChatErrorCode"; import { ChatLobbyType } from "./../src/chat/enums/ChatLobbyTypes"; -import { NotificationsService } from "./../src/notifications/notifications.service"; import { PruneDirectMessages } from "./../src/chat/jobs/PruneDirectMessages"; import { directRoomId } from "./../src/chat/utilities/directRoomId"; @@ -19,7 +18,12 @@ describe("direct messages (SQL-driven)", () => { let postgres: PostgresService; let fx: Fixtures; let chat: ChatService; - let bell: NotificationsService; + + const push = { + sendChatMessage: jest.fn(async () => {}), + retractChatMessage: jest.fn(async () => {}), + editChatMessage: jest.fn(async () => {}), + }; const logger = { log: jest.fn(), warn: jest.fn(), error: jest.fn() }; @@ -45,17 +49,6 @@ describe("direct messages (SQL-driven)", () => { postgres = db.postgres; fx = new Fixtures(postgres, 76561199400000000n); - bell = new NotificationsService( - {} as any, - postgres, - logger as any, - { get: () => ({ webDomain: "https://example.com" }) } as any, - {} as any, - {} as any, - {} as any, - {} as any, - ); - chat = new ChatService( logger as any, {} as any, @@ -69,17 +62,7 @@ describe("direct messages (SQL-driven)", () => { } as any, postgres, { getConnection: () => redis } as any, - { - notifyPlayers: jest.fn(), - markConversationRead: jest.fn(), - collapseOlderUnread: jest.fn(), - retractChatMessage: (messageId: string) => - bell.retractChatMessage(messageId), - retractChatMessageFromBlocked: (messageId: string) => - bell.retractChatMessageFromBlocked(messageId), - updateChatMessagePreview: (messageId: string, preview: string) => - bell.updateChatMessagePreview(messageId, preview), - } as any, + push as any, new PlayerBlocksService(postgres), ); }, 600_000); @@ -94,7 +77,6 @@ describe("direct messages (SQL-driven)", () => { await postgres.query("DELETE FROM direct_messages"); await postgres.query("DELETE FROM direct_conversations"); await postgres.query("DELETE FROM chat_read_state"); - await postgres.query("DELETE FROM notifications"); await postgres.query("DELETE FROM players"); }); @@ -449,27 +431,6 @@ describe("direct messages (SQL-driven)", () => { [id, minutes], ); - const bellRow = async (steamId: string, messageId: string) => { - const [row] = await postgres.query>( - `INSERT INTO notifications - (type, title, message, role, steam_id, entity_id, data) - VALUES ('ChatMessage', 'Someone', 'typo', 'user', $1::bigint, - 'direct:' || $1, jsonb_build_object('messageId', $2::text)) - RETURNING id::text AS id`, - [steamId, messageId], - ); - return row.id; - }; - - const notification = async (id: string) => - ( - await postgres.query< - Array<{ message: string; deleted_at: Date | null }> - >(`SELECT message, deleted_at FROM notifications WHERE id = $1::uuid`, [ - id, - ]) - ).at(0); - const as = (steamId: string) => ({ steam_id: steamId }) as any; it("stamps edited_at and leaves created_at alone", async () => { @@ -578,37 +539,29 @@ describe("direct messages (SQL-driven)", () => { ).resolves.toEqual({ edited: false, code: ChatErrorCode.NotFound }); }); - it("shows the edit on the recipient's unread bell row", async () => { + it("shows the edit on a push still being held", async () => { const me = await fx.player(); const friend = await fx.player(); const room = directRoomId(me, friend); const id = await sent(room, me); - const row = await bellRow(friend, id); await chat.editMessage(ChatLobbyType.Direct, room, id, as(me), "fixed"); - expect(await notification(row)).toEqual({ - message: "fixed", - deleted_at: null, - }); + expect(push.editChatMessage).toHaveBeenCalledWith(id, "fixed"); }); - it("deletes within the window and retracts the recipient's bell row", async () => { + it("deletes within the window and retracts a push still being held", async () => { const me = await fx.player(); const friend = await fx.player(); const room = directRoomId(me, friend); const id = await sent(room, me); - const row = await bellRow(friend, id); await expect( chat.deleteMessage(ChatLobbyType.Direct, room, id, as(me)), ).resolves.toEqual({ deleted: true }); expect(await stored(id)).toBeUndefined(); - expect(await notification(row)).toEqual({ - message: "", - deleted_at: expect.any(Date), - }); + expect(push.retractChatMessage).toHaveBeenCalledWith(id); const [{ count }] = await postgres.query>( `SELECT count(*)::text AS count FROM chat_message_deletions`, diff --git a/test/chat-moderation.spec.ts b/test/chat-moderation.spec.ts index 67b10cfd..58157ece 100644 --- a/test/chat-moderation.spec.ts +++ b/test/chat-moderation.spec.ts @@ -11,9 +11,12 @@ import { ChatService } from "./../src/chat/chat.service"; import { PlayerBlocksService } from "./../src/player-blocks/player-blocks.service"; import { ChatErrorCode } from "./../src/chat/enums/ChatErrorCode"; import { ChatLobbyType } from "./../src/chat/enums/ChatLobbyTypes"; -import { NotificationsService } from "./../src/notifications/notifications.service"; -import { NotificationPreferencesService } from "./../src/notifications/preferences/notification-preferences.service"; -import { PushNotificationsService } from "./../src/notifications/push/push-notifications.service"; +import { + ChatPush, + PushNotificationsService, +} from "./../src/notifications/push/push-notifications.service"; +import { chatThreadKey } from "./../src/notifications/push/notification-delivery"; +import { rolesAtOrAbove } from "./../src/utilities/isRoleAbove"; jest.mock("web-push", () => ({ setVapidDetails: jest.fn(), @@ -21,8 +24,8 @@ jest.mock("web-push", () => ({ generateVAPIDKeys: jest.fn(), })); -// The audit row, the redis removal and the bell retraction each live in a -// different store, and the unit specs stub all three. +// The audit row, the redis removal and the retraction of a held push each live +// in a different store, and the unit specs stub all three. describe("chat moderation (SQL-driven)", () => { let db: SqlTestDb; let postgres: PostgresService; @@ -30,7 +33,7 @@ describe("chat moderation (SQL-driven)", () => { let container: StartedTestContainer; let redis: Redis; let chat: ChatService; - let notifications: NotificationsService; + let push: PushNotificationsService; const logger = { log: jest.fn(), warn: jest.fn(), error: jest.fn() }; @@ -70,61 +73,6 @@ describe("chat moderation (SQL-driven)", () => { }), }); - const pushService = () => - new PushNotificationsService( - logger as any, - postgres, - { - get: (key: string) => - key === "app" - ? { webDomain: "https://example.com" } - : { - publicKey: "public-key", - privateKey: "private-key", - subject: "https://example.com", - }, - } as any, - { add: async () => ({}) } as any, - { - getConnection: () => ({ - exists: async () => 0, - set: async () => "OK", - get: async (): Promise => null, - ttl: async () => -2, - del: async () => 1, - rpush: async () => 1, - expire: async () => 1, - multi: () => ({ - lrange() { - return this; - }, - del() { - return this; - }, - rpush() { - return this; - }, - expire() { - return this; - }, - exec: async (): Promise> => [[null, []]], - }), - pipeline: () => { - const queued: string[] = []; - return { - set: () => {}, - hvals: (key: string) => queued.push(key), - exec: async (): Promise> => - queued.map(() => [null, []] as [unknown, Array]), - }; - }, - subscribe: async () => 1, - publish: async () => 1, - on: () => {}, - }), - } as any, - ); - beforeAll(async () => { container = await new GenericContainer("redis:8.8-alpine") .withExposedPorts(6379) @@ -138,16 +86,23 @@ describe("chat moderation (SQL-driven)", () => { postgres = db.postgres; fx = new Fixtures(postgres, 76561199600000000n); - notifications = new NotificationsService( - hasura() as any, - postgres, + push = new PushNotificationsService( logger as any, - { get: () => ({ webDomain: "https://example.com" }) } as any, - new NotificationPreferencesService(postgres), - { add: jest.fn() } as any, - { add: jest.fn() } as any, - { add: jest.fn() } as any, + postgres, + { + get: (key: string) => + key === "app" + ? { webDomain: "https://example.com" } + : { + publicKey: "public-key", + privateKey: "private-key", + subject: "https://example.com", + }, + } as any, + { add: async () => ({}) } as any, + { getConnection: () => redis } as any, ); + await push.loadKeys(); chat = new ChatService( logger as any, @@ -155,7 +110,7 @@ describe("chat moderation (SQL-driven)", () => { hasura() as any, postgres, { getConnection: () => redis } as any, - notifications, + push, new PlayerBlocksService(postgres), ); }, 600_000); @@ -172,7 +127,8 @@ describe("chat moderation (SQL-driven)", () => { await postgres.query("DELETE FROM chat_message_deletions"); await postgres.query("DELETE FROM player_sanctions"); await postgres.query("DELETE FROM push_subscriptions"); - await postgres.query("DELETE FROM notifications"); + await postgres.query("DELETE FROM chat_read_state"); + await postgres.query("DELETE FROM player_blocks"); await postgres.query("DELETE FROM players"); }); @@ -312,20 +268,6 @@ describe("chat moderation (SQL-driven)", () => { expect(await audits()).toHaveLength(1); }); - it("tells a deleted message from one that merely expired", async () => { - const mod = await moderator(); - const author = await fx.player("Author"); - const matchId = randomUUID(); - const deleted = await post(matchId, author); - const expired = await post(matchId, author, "fine"); - - await chat.deleteMessage(ChatLobbyType.Match, matchId, deleted, mod); - await redis.hdel(`chat_match_${matchId}`, expired); - - expect(await chat["wasDeleted"](deleted)).toBe(true); - expect(await chat["wasDeleted"](expired)).toBe(false); - }); - it("keeps the audit when the author's player row goes", async () => { const mod = await moderator(); const author = await fx.player("Author"); @@ -344,30 +286,17 @@ describe("chat moderation (SQL-driven)", () => { }); }); - describe("retracting the bell", () => { - const chatNotification = async ( - steamId: string, - entityId: string, - messageId: string, - type = "MatchChatMessage", - ) => { - const [row] = await postgres.query>( - `INSERT INTO notifications - (type, title, message, role, steam_id, entity_id, data) - VALUES ($4, 'Author', 'something awful', 'user', - $1::bigint, $2, - jsonb_build_object('threadKey', 'chat:' || $2, - 'messageId', $3::text)) - RETURNING id::text AS id`, - [steamId, entityId, messageId, type], - ); - return row.id; - }; + describe("a push still being held", () => { + const pushed = () => (webPush.sendNotification as jest.Mock).mock.calls; - // Match chat is off for push by default, so delivery is shown on a room - // whose category is on. - const subscribedReader = async () => { + const bodyOf = (call: number) => JSON.parse(pushed()[call][1]).body; + + const subscribedReader = async (role = "user") => { const reader = await fx.player("Reader"); + await postgres.query( + `UPDATE players SET role = $2 WHERE steam_id = $1::bigint`, + [reader, role], + ); await postgres.query( `INSERT INTO push_subscriptions (steam_id, endpoint, p256dh, auth) VALUES ($1::bigint, $2, 'key', 'auth')`, @@ -376,213 +305,241 @@ describe("chat moderation (SQL-driven)", () => { return reader; }; - const deliver = async (id: string) => { - const push = pushService(); - await push.loadKeys(); - await push.sendForNotification({ id, type: "ChatMessage" }); - }; + // Plain chat rather than match chat, which is off for push by default. + const message = ( + matchId: string, + messageId: string, + senderSteamId: string, + ): ChatPush => ({ + messageId, + type: "ChatMessage", + title: "Author", + message: "something awful", + entityId: `match:${matchId}`, + threadKey: chatThreadKey(ChatLobbyType.Match, matchId), + threadLabel: "Blue vs Red", + senderSteamId, + blockExemptRoles: rolesAtOrAbove("moderator"), + }); - const notification = async (id: string) => - ( - await postgres.query< - Array<{ deleted_at: Date | null; message: string }> - >(`SELECT deleted_at, message FROM notifications WHERE id = $1::uuid`, [ - id, - ]) - ).at(0); + // The first message buzzes; the rest wait for the window to close. + const burst = async ( + reader: string, + matchId: string, + sender: string, + messageIds: string[], + ) => { + for (const messageId of messageIds) { + await push.sendChatMessage( + [reader], + message(matchId, messageId, sender), + ); + } - const deletedAt = async (id: string) => - (await notification(id))?.deleted_at; + expect(pushed()).toHaveLength(1); + }; - it("retracts the deleted message's row and no other", async () => { + const closeWindow = (reader: string, matchId: string) => + push.sendPending(reader, chatThreadKey(ChatLobbyType.Match, matchId)); + + it("leaves a message a moderator deleted out of the summary", async () => { const mod = await moderator(); const author = await fx.player("Author"); - const reader = await fx.player("Reader"); + const reader = await subscribedReader(); const matchId = randomUUID(); + const first = await post(matchId, author); const target = await post(matchId, author); - const other = await post(matchId, author, "fine"); - - const retracted = await chatNotification( - reader, - `match:${matchId}`, - target, - ); - const kept = await chatNotification(reader, `match:${matchId}`, other); + const last = await post(matchId, author); + await burst(reader, matchId, author, [first, target, last]); await chat.deleteMessage(ChatLobbyType.Match, matchId, target, mod); + await closeWindow(reader, matchId); - expect(await deletedAt(retracted)).toBeInstanceOf(Date); - expect(await deletedAt(kept)).toBeNull(); + expect(pushed()).toHaveLength(2); + expect(bodyOf(1)).toBe("2 new messages"); }); - it("takes the text out of the row, so it cannot be read back", async () => { + it("says nothing more when the only message after the first was deleted", async () => { const mod = await moderator(); const author = await fx.player("Author"); - const reader = await fx.player("Reader"); + const reader = await subscribedReader(); const matchId = randomUUID(); + const first = await post(matchId, author); const target = await post(matchId, author); - const id = await chatNotification(reader, `match:${matchId}`, target); + await burst(reader, matchId, author, [first, target]); await chat.deleteMessage(ChatLobbyType.Match, matchId, target, mod); + await closeWindow(reader, matchId); - expect((await notification(id))?.message).toBe(""); + expect(pushed()).toHaveLength(1); }); - it("blanks a row the bell had already collapsed, leaving it retired", async () => { - const reader = await fx.player("Reader"); - const messageId = randomUUID(); - const id = await chatNotification(reader, "tournament:t-1", messageId); - await postgres.query( - `UPDATE notifications - SET deleted_at = now() - interval '1 hour' - WHERE id = $1::uuid`, - [id], - ); - const before = await deletedAt(id); + it("does not push a message deleted before its push went out", async () => { + const mod = await moderator(); + const author = await fx.player("Author"); + const reader = await subscribedReader(); + const matchId = randomUUID(); + const target = await post(matchId, author); - await notifications.retractChatMessage(messageId); + await chat.deleteMessage(ChatLobbyType.Match, matchId, target, mod); + await push.sendChatMessage([reader], message(matchId, target, author)); - expect(await notification(id)).toEqual({ - deleted_at: before, - message: "", - }); + expect(pushed()).toHaveLength(0); }); it("retracts a draft lobby's message after it moved into the match", async () => { const mod = await moderator(); const author = await fx.player("Author"); - const reader = await fx.player("Reader"); + const reader = await subscribedReader(); const draftId = randomUUID(); const matchId = randomUUID(); + const first = await post(draftId, author); const target = await post(draftId, author); await redis.rename(`chat_match_${draftId}`, `chat_draft_${draftId}`); - const id = await chatNotification( - reader, - `draft:${draftId}`, - target, - "ChatMessage", - ); - + await burst(reader, matchId, author, [first, target]); await chat.migrateLobbyMessages( ChatLobbyType.Draft, draftId, ChatLobbyType.Match, matchId, ); - await expect( chat.deleteMessage(ChatLobbyType.Match, matchId, target, mod), ).resolves.toEqual({ deleted: true }); + await closeWindow(reader, matchId); - expect(await deletedAt(id)).toBeInstanceOf(Date); + expect(pushed()).toHaveLength(1); }); - it("leaves a message's rows alone when it expired rather than being deleted", async () => { - // A 0 TTL drops the field as soon as it is written, and a draft lobby's - // history moves out from under it into the match. Neither is a delete. + it("shows what an edit made of a held message", async () => { const author = await fx.player("Author"); - const reader = await fx.player("Reader"); + const reader = await subscribedReader(); const matchId = randomUUID(); - const messageId = await post(matchId, author); - const id = await chatNotification(reader, `match:${matchId}`, messageId); - await redis.hdel(`chat_match_${matchId}`, messageId); - - const members = jest - .spyOn(chat, "getLobbyMemberSteamIds") - .mockResolvedValueOnce([author, reader]); - const written = jest - .spyOn(notifications, "notifyPlayers") - .mockResolvedValueOnce(undefined); - - await chat["notifyLobbyMembers"]( - ChatLobbyType.Match, - matchId, - { steam_id: author, name: "Author", role: "user" } as any, - "Author", - "something awful", - messageId, + const target = randomUUID(); + + await postgres.query( + `UPDATE players + SET quiet_hours_start = (now() AT TIME ZONE 'UTC')::time - interval '1 hour', + quiet_hours_end = (now() AT TIME ZONE 'UTC')::time + interval '1 hour', + notification_timezone = 'UTC' + WHERE steam_id = $1::bigint`, + [reader], + ); + await push.sendChatMessage([reader], message(matchId, target, author)); + await postgres.query( + `UPDATE players SET quiet_hours_start = NULL, quiet_hours_end = NULL + WHERE steam_id = $1::bigint`, + [reader], ); - members.mockRestore(); - written.mockRestore(); + expect(pushed()).toHaveLength(0); - expect(await notification(id)).toEqual({ - deleted_at: null, - message: "something awful", - }); + await push.editChatMessage(target, "fixed"); + await closeWindow(reader, matchId); + + expect(pushed()).toHaveLength(1); + expect(bodyOf(0)).toBe("fixed"); }); - it("leaves a notification that is not chat alone, whatever its data holds", async () => { - const reader = await fx.player("Reader"); - const messageId = randomUUID(); - const id = await chatNotification( - reader, - "tournament:t-1", - messageId, - "MatchStatusChange", + it("drops what a sender the reader has since blocked said", async () => { + const author = await fx.player("Author"); + const reader = await subscribedReader(); + const matchId = randomUUID(); + + await burst(reader, matchId, author, [randomUUID(), randomUUID()]); + await postgres.query( + `INSERT INTO player_blocks (blocker_steam_id, blocked_steam_id) + VALUES ($1::bigint, $2::bigint)`, + [reader, author], ); + await closeWindow(reader, matchId); - await notifications.retractChatMessage(messageId); + expect(pushed()).toHaveLength(1); + }); - expect(await notification(id)).toEqual({ - deleted_at: null, - message: "something awful", - }); + it("still tells a moderator who blocked the sender about a group room", async () => { + const author = await fx.player("Author"); + const reader = await subscribedReader("moderator"); + const matchId = randomUUID(); + + await burst(reader, matchId, author, [randomUUID(), randomUUID()]); + await postgres.query( + `INSERT INTO player_blocks (blocker_steam_id, blocked_steam_id) + VALUES ($1::bigint, $2::bigint)`, + [reader, author], + ); + await closeWindow(reader, matchId); + + expect(pushed()).toHaveLength(2); + expect(bodyOf(1)).toBe("2 new messages"); }); - it("finds the message's rows through an index", async () => { - const plan = await postgres.transaction(async (client) => { - await client.query("SET LOCAL enable_seqscan = off"); + it("says nothing about held messages the reader has since read", async () => { + const author = await fx.player("Author"); + const reader = await subscribedReader(); + const matchId = randomUUID(); - const { rows } = await client.query( - `EXPLAIN SELECT id FROM notifications - WHERE data->>'messageId' = $1`, - [randomUUID()], - ); + await burst(reader, matchId, author, [randomUUID(), randomUUID()]); + await postgres.query( + `INSERT INTO chat_read_state (steam_id, thread, last_read_at) + VALUES ($1::bigint, $2, now())`, + [reader, chatThreadKey(ChatLobbyType.Match, matchId)], + ); + await closeWindow(reader, matchId); - return rows.map((row) => row["QUERY PLAN"]).join("\n"); - }); + expect(pushed()).toHaveLength(1); + }); + + it("ignores a cursor on a different thread", async () => { + const author = await fx.player("Author"); + const reader = await subscribedReader(); + const matchId = randomUUID(); + + await burst(reader, matchId, author, [randomUUID(), randomUUID()]); + await postgres.query( + `INSERT INTO chat_read_state (steam_id, thread, last_read_at) + VALUES ($1::bigint, $2, now())`, + [reader, chatThreadKey(ChatLobbyType.Match, randomUUID())], + ); + await closeWindow(reader, matchId); - expect(plan).toContain("notifications_message_id_idx"); + expect(pushed()).toHaveLength(2); }); - it("drops a retracted row from push delivery", async () => { + it("pushes nobody who turned chat off", async () => { + const author = await fx.player("Author"); const reader = await subscribedReader(); - const messageId = randomUUID(); - const id = await chatNotification( - reader, - "tournament:t-1", - messageId, - "ChatMessage", + await postgres.query( + `INSERT INTO notification_preferences (steam_id, channel, key, enabled) + VALUES ($1::bigint, 'push', 'chat', false)`, + [reader], ); - await notifications.retractChatMessage(messageId); - await deliver(id); + await push.sendChatMessage( + [reader], + message(randomUUID(), randomUUID(), author), + ); - expect(webPush.sendNotification).not.toHaveBeenCalled(); + expect(pushed()).toHaveLength(0); }); - it("still delivers the row next to it", async () => { + it("writes no notifications row", async () => { + const author = await fx.player("Author"); const reader = await subscribedReader(); - const retracted = randomUUID(); - await chatNotification( - reader, - "tournament:t-1", - retracted, - "ChatMessage", - ); - const kept = await chatNotification( - reader, - "tournament:t-1", - randomUUID(), - "ChatMessage", + + await push.sendChatMessage( + [reader], + message(randomUUID(), randomUUID(), author), ); - await notifications.retractChatMessage(retracted); - await deliver(kept); + const [{ count }] = await postgres.query>( + `SELECT count(*)::text AS count FROM notifications + WHERE type IN ('ChatMessage', 'MatchChatMessage')`, + ); - expect(webPush.sendNotification).toHaveBeenCalledTimes(1); + expect(pushed()).toHaveLength(1); + expect(count).toBe("0"); }); }); diff --git a/test/chat-redis-actions.spec.ts b/test/chat-redis-actions.spec.ts index 5a19f677..5e7b86cf 100644 --- a/test/chat-redis-actions.spec.ts +++ b/test/chat-redis-actions.spec.ts @@ -11,8 +11,6 @@ import { ChatGateway } from "./../src/chat/chat.gateway"; import { PlayerBlocksService } from "./../src/player-blocks/player-blocks.service"; import { ChatErrorCode } from "./../src/chat/enums/ChatErrorCode"; import { ChatLobbyType } from "./../src/chat/enums/ChatLobbyTypes"; -import { NotificationsService } from "./../src/notifications/notifications.service"; -import { NotificationPreferencesService } from "./../src/notifications/preferences/notification-preferences.service"; // The edit is one compare-and-set script against real hash-field expiry, which // no fake reproduces: HSET dropping a field's TTL is the whole reason it exists. @@ -23,7 +21,12 @@ describe("chat edits and self deletes (SQL-driven)", () => { let container: StartedTestContainer; let redis: Redis; let chat: ChatService; - let notifications: NotificationsService; + + const push = { + sendChatMessage: jest.fn(async () => {}), + retractChatMessage: jest.fn(async () => {}), + editChatMessage: jest.fn(async () => {}), + }; const logger = { log: jest.fn(), warn: jest.fn(), error: jest.fn() }; @@ -72,24 +75,13 @@ describe("chat edits and self deletes (SQL-driven)", () => { postgres = db.postgres; fx = new Fixtures(postgres, 76561199610000000n); - notifications = new NotificationsService( - hasura() as any, - postgres, - logger as any, - { get: () => ({ webDomain: "https://example.com" }) } as any, - new NotificationPreferencesService(postgres), - { add: jest.fn() } as any, - { add: jest.fn() } as any, - { add: jest.fn() } as any, - ); - chat = new ChatService( logger as any, {} as any, hasura() as any, postgres, { getConnection: () => redis } as any, - notifications, + push as any, new PlayerBlocksService(postgres), ); }, 600_000); @@ -107,7 +99,6 @@ describe("chat edits and self deletes (SQL-driven)", () => { await postgres.query("DELETE FROM chat_message_deletions"); await postgres.query("DELETE FROM chat_message_edits"); await postgres.query("DELETE FROM player_sanctions"); - await postgres.query("DELETE FROM notifications"); await postgres.query("DELETE FROM players"); }); @@ -1103,126 +1094,6 @@ describe("chat edits and self deletes (SQL-driven)", () => { }); }); - describe("the bell's preview", () => { - const row = async ( - steamId: string, - messageId: string, - overrides: { - is_read?: boolean; - deleted?: boolean; - message?: string; - } = {}, - ) => { - const [inserted] = await postgres.query>( - `INSERT INTO notifications - (type, title, message, role, steam_id, entity_id, data, - is_read, deleted_at) - VALUES ('MatchChatMessage', 'Author', $3, 'user', $1::bigint, - 'match:m-1', - jsonb_build_object('messageId', $2::text), - $4, CASE WHEN $5 THEN now() END) - RETURNING id::text AS id`, - [ - steamId, - messageId, - overrides.message ?? "typo", - overrides.is_read ?? false, - overrides.deleted ?? false, - ], - ); - return inserted.id; - }; - - const text = async (id: string) => - ( - await postgres.query>( - `SELECT message FROM notifications WHERE id = $1::uuid`, - [id], - ) - )[0].message; - - it("shows the edited text on the recipient's unread row", async () => { - const user = await author(); - const reader = await fx.player("Reader"); - const matchId = randomUUID(); - const id = await place(matchId, user); - const unread = await row(reader, id); - - await edit(matchId, id, user, "fixed"); - - expect(await text(unread)).toBe("<b>fixed</b>"); - }); - - // A read row stays in the bell, and a collapsed one can be restored by - // its recipient, so neither may keep the text the author took back. - it("rewrites every row for that message, read or collapsed", async () => { - const reader = await fx.player("Reader"); - const messageId = randomUUID(); - const unread = await row(reader, messageId); - const read = await row(reader, messageId, { is_read: true }); - const collapsed = await row(reader, messageId, { deleted: true }); - const other = await row(reader, randomUUID()); - - await notifications.updateChatMessagePreview(messageId, "fixed"); - - expect(await text(unread)).toBe("fixed"); - expect(await text(read)).toBe("fixed"); - expect(await text(collapsed)).toBe("fixed"); - expect(await text(other)).toBe("typo"); - }); - - it("leaves a row retracted mid-edit blank", async () => { - const reader = await fx.player("Reader"); - const messageId = randomUUID(); - const id = await row(reader, messageId); - let pending: Promise = Promise.resolve(); - - await postgres.transaction(async (client) => { - await client.query( - `UPDATE notifications SET deleted_at = now(), message = '' - WHERE id = $1::uuid`, - [id], - ); - - pending = notifications.updateChatMessagePreview(messageId, "fixed"); - await new Promise((resolve) => setTimeout(resolve, 200)); - }); - - await pending; - - expect(await text(id)).toBe(""); - }); - - it("leaves a retracted row blank when the edit lands after the delete", async () => { - const reader = await fx.player("Reader"); - const messageId = randomUUID(); - const id = await row(reader, messageId); - - await notifications.retractChatMessage(messageId); - await notifications.updateChatMessagePreview(messageId, "fixed"); - - expect(await text(id)).toBe(""); - }); - - it("finds the message's rows through an index", async () => { - const plan = await postgres.transaction(async (client) => { - await client.query("SET LOCAL enable_seqscan = off"); - - const { rows } = await client.query( - `EXPLAIN UPDATE notifications SET message = 'x' - WHERE data->>'messageId' = $1 - AND type IN ('ChatMessage', 'MatchChatMessage') - AND message <> ''`, - [randomUUID()], - ); - - return rows.map((plan) => plan["QUERY PLAN"]).join("\n"); - }); - - expect(plan).toContain("notifications_message_id_idx"); - }); - }); - describe("sending from the web", () => { const LIMIT = 5; const WINDOW_MS = 3_000; @@ -1259,7 +1130,7 @@ describe("chat edits and self deletes (SQL-driven)", () => { stub as any, postgres, { getConnection: () => redis } as any, - notifications, + push as any, new PlayerBlocksService(postgres), ); diff --git a/test/drop-chat-notifications.spec.ts b/test/drop-chat-notifications.spec.ts new file mode 100644 index 00000000..793d5d9d --- /dev/null +++ b/test/drop-chat-notifications.spec.ts @@ -0,0 +1,99 @@ +import { readFileSync } from "fs"; +import { join } from "path"; +import { bootMigratedDb, SqlTestDb } from "./utils/sql-test-db"; + +// Chat pushes straight from the conversation now, so the rows it used to write +// would only sit in the bell, and an in-app chat preference has no toggle left +// to change it back. +describe("drop chat notifications migration", () => { + let db: SqlTestDb; + + const up = readFileSync( + join( + __dirname, + "../hasura/migrations/default/1889000001600_drop_chat_notifications/up.sql", + ), + "utf8", + ); + + const STEAM_ID = "76561199500000021"; + + const typesLeft = async () => + ( + await db.postgres.query>( + `SELECT type FROM public.notifications + WHERE steam_id = $1 + ORDER BY type`, + [STEAM_ID], + ) + ).map(({ type }) => type); + + const preferenceKeys = async () => + ( + await db.postgres.query>( + `SELECT channel, key FROM public.notification_preferences + WHERE steam_id = $1 + ORDER BY channel, key`, + [STEAM_ID], + ) + ).map(({ channel, key }) => `${channel}:${key}`); + + beforeAll(async () => { + db = await bootMigratedDb("DropChatNotificationsTest"); + + await db.postgres.query( + `INSERT INTO public.players (steam_id, name) VALUES ($1, 'Chatty') + ON CONFLICT (steam_id) DO NOTHING`, + [STEAM_ID], + ); + + for (const type of ["ChatMessage", "MatchChatMessage", "MatchImported"]) { + await db.postgres.query( + `INSERT INTO public.notifications + (type, title, message, role, steam_id, entity_id) + VALUES ($2, 'DrClampz', 'gg', 'user', $1, 'direct:1:2')`, + [STEAM_ID, type], + ); + } + + for (const [channel, key] of [ + ["in_app", "ChatMessage"], + ["in_app", "MatchChatMessage"], + ["in_app", "MatchImported"], + ["push", "chat"], + ]) { + await db.postgres.query( + `INSERT INTO public.notification_preferences + (steam_id, channel, key, enabled) + VALUES ($1, $2, $3, false)`, + [STEAM_ID, channel, key], + ); + } + + await db.postgres.query(up); + }, 600_000); + + afterAll(async () => { + await db?.stop(); + }); + + it("deletes the chat rows and nothing else", async () => { + expect(await typesLeft()).toEqual(["MatchImported"]); + }); + + it("drops the index that only found chat rows", async () => { + const [{ count }] = await db.postgres.query>( + `SELECT count(*)::text AS count FROM pg_indexes + WHERE indexname = 'notifications_message_id_idx'`, + ); + + expect(count).toBe("0"); + }); + + it("drops only the in-app chat preferences", async () => { + expect(await preferenceKeys()).toEqual([ + "in_app:MatchImported", + "push:chat", + ]); + }); +}); diff --git a/test/notifications.spec.ts b/test/notifications.spec.ts index 5e973846..874fa233 100644 --- a/test/notifications.spec.ts +++ b/test/notifications.spec.ts @@ -172,10 +172,13 @@ describe("notifications (SQL-driven)", () => { const listening = await fx.player(); const service = preferences(); - await service.set(muted, "in_app", "ChatMessage", false); + await service.set(muted, "in_app", "MatchImported", false); expect( - await service.filterInAppRecipients("ChatMessage", [muted, listening]), + await service.filterInAppRecipients("MatchImported", [ + muted, + listening, + ]), ).toEqual([listening]); }); @@ -425,58 +428,6 @@ describe("notifications (SQL-driven)", () => { }); }); - describe("collapseOlderUnread", () => { - it("keeps only the newest unread row for a conversation", async () => { - const steamId = await fx.player(); - - for (const body of ["first", "second", "third"]) { - await postgres.query( - `INSERT INTO notifications (type, title, message, role, steam_id, entity_id) - VALUES ('ChatMessage', 'Someone', $1, 'user', $2::bigint, 'match:m-1')`, - [body, steamId], - ); - } - - await notifications().collapseOlderUnread("ChatMessage", "match:m-1", [ - steamId, - ]); - - const rows = await postgres.query>( - `SELECT message FROM notifications - WHERE deleted_at IS NULL AND type = 'ChatMessage'`, - ); - - expect(rows.map((row) => row.message)).toEqual(["third"]); - }); - - it("leaves another conversation untouched", async () => { - const steamId = await fx.player(); - - for (const entity of ["match:m-1", "match:m-2"]) { - await postgres.query( - `INSERT INTO notifications (type, title, message, role, steam_id, entity_id) - VALUES ('ChatMessage', 'Someone', 'hi', 'user', $1::bigint, $2)`, - [steamId, entity], - ); - } - - await notifications().collapseOlderUnread("ChatMessage", "match:m-1", [ - steamId, - ]); - - const [row] = await postgres.query>( - `SELECT count(*)::text AS count FROM notifications WHERE deleted_at IS NULL`, - ); - expect(row.count).toBe("2"); - }); - }); - - // A notification that reaches the support webhook is posted verbatim into a - // staff channel, so only the types an operator has to act on go there and - // everything else is in-app. Routing used to say that the other way round -- - // a list of exclusions, with Discord the default -- which is how invites and - // check-in reminders addressed to one player came to be posted to staff. - // // The next type that gets routed through notifyPlayers, which is where the // webhook lives, should fail here rather than in a staff channel. describe("discord relay", () => { @@ -594,8 +545,6 @@ describe("notifications (SQL-driven)", () => { }); describe("push recipient resolution", () => { - // What a bundling window has waiting in it, for the trailing-summary case. - let pending: string[] = []; // Jobs the delivery queue was handed, so a deferred push can be told apart // from a dropped one. let pendingQueued: Array<{ delay?: number }> = []; @@ -624,7 +573,7 @@ describe("notifications (SQL-driven)", () => { expire() { return this; }, - exec: async (): Promise> => [[null, pending]], + exec: async (): Promise> => [[null, []]], }), pipeline: () => { const queued: string[] = []; @@ -688,19 +637,17 @@ describe("notifications (SQL-driven)", () => { return push; }; - const chatNotification = async ( + const matchAlert = async ( steamId: string, - overrides: { createdAt?: string; isRead?: boolean } = {}, + overrides: { isRead?: boolean } = {}, ) => { const [row] = await postgres.query>( `INSERT INTO notifications - (type, title, message, role, steam_id, entity_id, is_read, data, created_at) - VALUES ('ChatMessage', 'Luke', 'hey', 'user', $1::bigint, - 'match:m-1', $2, - '{"threadKey":"chat:match:m-1"}'::jsonb, - COALESCE($3::timestamptz, now())) + (type, title, message, role, steam_id, entity_id, is_read) + VALUES ('MatchStatusChange', 'Match paused', 'paused', 'user', + $1::bigint, 'm-1', $2) RETURNING id::text AS id`, - [steamId, overrides.isRead ?? false, overrides.createdAt ?? null], + [steamId, overrides.isRead ?? false], ); return row.id; }; @@ -721,73 +668,11 @@ describe("notifications (SQL-driven)", () => { it("resolves a subscribed recipient", async () => { const steamId = await fx.player(); await subscribe(steamId); - const id = await chatNotification(steamId); - - await (await configuredService()).sendForNotification({ - id, - type: "ChatMessage", - }); - - expect(webPush.sendNotification).toHaveBeenCalledTimes(1); - }); - - it("says nothing when the thread was read after the message", async () => { - // The whole point of the read cursor: a message the recipient has - // already scrolled past is not worth a buzz. - const steamId = await fx.player(); - await subscribe(steamId); - const id = await chatNotification(steamId, { - createdAt: new Date(Date.now() - 60_000).toISOString(), - }); - - await postgres.query( - `INSERT INTO chat_read_state (steam_id, thread, last_read_at) - VALUES ($1::bigint, 'chat:match:m-1', now())`, - [steamId], - ); - - await (await configuredService()).sendForNotification({ - id, - type: "ChatMessage", - }); - - expect(webPush.sendNotification).not.toHaveBeenCalled(); - }); - - it("still buzzes for a message newer than the cursor", async () => { - const steamId = await fx.player(); - await subscribe(steamId); - - await postgres.query( - `INSERT INTO chat_read_state (steam_id, thread, last_read_at) - VALUES ($1::bigint, 'chat:match:m-1', now() - interval '1 hour')`, - [steamId], - ); - - const id = await chatNotification(steamId); - - await (await configuredService()).sendForNotification({ - id, - type: "ChatMessage", - }); - - expect(webPush.sendNotification).toHaveBeenCalledTimes(1); - }); - - it("ignores a cursor for a different thread", async () => { - const steamId = await fx.player(); - await subscribe(steamId); - const id = await chatNotification(steamId); - - await postgres.query( - `INSERT INTO chat_read_state (steam_id, thread, last_read_at) - VALUES ($1::bigint, 'chat:match:m-2', now())`, - [steamId], - ); + const id = await matchAlert(steamId); await (await configuredService()).sendForNotification({ id, - type: "ChatMessage", + type: "MatchStatusChange", }); expect(webPush.sendNotification).toHaveBeenCalledTimes(1); @@ -836,11 +721,11 @@ describe("notifications (SQL-driven)", () => { it("drops a row already dealt with in the bell", async () => { const steamId = await fx.player(); await subscribe(steamId); - const id = await chatNotification(steamId, { isRead: true }); + const id = await matchAlert(steamId, { isRead: true }); await (await configuredService()).sendForNotification({ id, - type: "ChatMessage", + type: "MatchStatusChange", }); expect(webPush.sendNotification).not.toHaveBeenCalled(); @@ -872,55 +757,16 @@ describe("notifications (SQL-driven)", () => { // which is what makes the flag above load-bearing rather than decorative. const steamId = await fx.player(); await subscribe(steamId); - const id = await chatNotification(steamId, { isRead: true }); + const id = await matchAlert(steamId, { isRead: true }); await (await configuredService()).sendForNotification({ id, - type: "ChatMessage", + type: "MatchStatusChange", }); expect(webPush.sendNotification).not.toHaveBeenCalled(); }); - it("still counts a burst whose older rows the bell collapsed", async () => { - // collapseOlderUnread soft-deletes every superseded ChatMessage row so - // the bell shows one entry per conversation. The summary's count comes - // from the window rather than from those rows for exactly that reason -- - // resolving them finds one survivor and would report a burst of three as - // a single message. - const steamId = await fx.player(); - await subscribe(steamId); - - const ids = [ - await chatNotification(steamId, { - createdAt: new Date(Date.now() - 3000).toISOString(), - }), - await chatNotification(steamId, { - createdAt: new Date(Date.now() - 2000).toISOString(), - }), - await chatNotification(steamId), - ]; - - await notifications().collapseOlderUnread("ChatMessage", "match:m-1", [ - steamId, - ]); - - const [surviving] = await postgres.query>( - `SELECT count(*)::text AS count FROM notifications - WHERE type = 'ChatMessage' AND deleted_at IS NULL`, - ); - expect(surviving.count).toBe("1"); - - pending = ids; - await (await configuredService()).sendPending(steamId, "chat:match:m-1"); - - const [, payload] = (webPush.sendNotification as jest.Mock).mock.calls[0]; - expect(JSON.parse(payload)).toMatchObject({ - body: "3 new messages", - count: 3, - }); - }); - it("holds rather than drops during the recipient's quiet hours", async () => { const steamId = await fx.player(); await subscribe(steamId); @@ -932,12 +778,12 @@ describe("notifications (SQL-driven)", () => { WHERE steam_id = $1::bigint`, [steamId], ); - const id = await chatNotification(steamId); + const id = await matchAlert(steamId); pendingQueued = []; await (await configuredService()).sendForNotification({ id, - type: "ChatMessage", + type: "MatchStatusChange", }); expect(webPush.sendNotification).not.toHaveBeenCalled();