diff --git a/hasura/metadata/databases/default/tables/public_direct_message_reactions.yaml b/hasura/metadata/databases/default/tables/public_direct_message_reactions.yaml new file mode 100644 index 00000000..4bd37572 --- /dev/null +++ b/hasura/metadata/databases/default/tables/public_direct_message_reactions.yaml @@ -0,0 +1,3 @@ +table: + name: direct_message_reactions + schema: public diff --git a/hasura/metadata/databases/default/tables/tables.yaml b/hasura/metadata/databases/default/tables/tables.yaml index 07d8efd5..e90a8ce8 100644 --- a/hasura/metadata/databases/default/tables/tables.yaml +++ b/hasura/metadata/databases/default/tables/tables.yaml @@ -12,6 +12,7 @@ - "!include public_custom_pages.yaml" - "!include public_db_backups.yaml" - "!include public_direct_conversations.yaml" +- "!include public_direct_message_reactions.yaml" - "!include public_direct_messages.yaml" - "!include public_draft_game_picks.yaml" - "!include public_draft_game_players.yaml" diff --git a/hasura/migrations/default/1888000000300_direct_message_reactions/down.sql b/hasura/migrations/default/1888000000300_direct_message_reactions/down.sql new file mode 100644 index 00000000..073f9f5a --- /dev/null +++ b/hasura/migrations/default/1888000000300_direct_message_reactions/down.sql @@ -0,0 +1 @@ +DROP TABLE IF EXISTS public.direct_message_reactions; diff --git a/hasura/migrations/default/1888000000300_direct_message_reactions/up.sql b/hasura/migrations/default/1888000000300_direct_message_reactions/up.sql new file mode 100644 index 00000000..8cc66450 --- /dev/null +++ b/hasura/migrations/default/1888000000300_direct_message_reactions/up.sql @@ -0,0 +1,13 @@ +CREATE TABLE IF NOT EXISTS public.direct_message_reactions ( + message_id uuid NOT NULL REFERENCES public.direct_messages (id) + ON DELETE CASCADE, + steam_id bigint NOT NULL REFERENCES public.players (steam_id) + ON UPDATE CASCADE ON DELETE CASCADE, + reaction text NOT NULL, + created_at timestamptz NOT NULL DEFAULT now(), + + PRIMARY KEY (message_id, steam_id, reaction) +); + +CREATE INDEX IF NOT EXISTS direct_message_reactions_steam_id_idx + ON public.direct_message_reactions (steam_id); diff --git a/src/chat/chat.gateway.spec.ts b/src/chat/chat.gateway.spec.ts index 4121b04b..c2dd274b 100644 --- a/src/chat/chat.gateway.spec.ts +++ b/src/chat/chat.gateway.spec.ts @@ -637,3 +637,150 @@ describe("ChatGateway lobby:edit", () => { expect(chat.sendChatToServer).not.toHaveBeenCalled(); }); }); + +describe("ChatGateway lobby:react", () => { + const MESSAGE_ID = "3f0c1d2e-4b5a-4c6d-8e7f-9a0b1c2d3e4f"; + + let chat: { toggleReaction: jest.Mock; sendChatToServer: jest.Mock }; + let gateway: ChatGateway; + + const client = (user: any = { steam_id: "1", name: "Luke", role: "user" }) => + ({ id: "client-1", user, send: jest.fn() }) as any; + + const sent = (socket: { send: jest.Mock }) => + socket.send.mock.calls.map(([raw]) => JSON.parse(raw)); + + const reaction = (overrides: Record = {}) => ({ + id: "m-1", + type: ChatLobbyType.Match, + messageId: MESSAGE_ID, + reaction: "heart", + ...overrides, + }); + + beforeEach(() => { + chat = { + toggleReaction: jest.fn().mockResolvedValue({ + toggled: true, + reactions: { heart: ["1"] }, + }), + sendChatToServer: jest.fn(), + }; + gateway = new ChatGateway(chat as any); + }); + + it("ignores a socket that has not signed in", async () => { + const socket = client(null); + + await gateway.react(reaction({ requestId: "r-1" }) as any, socket); + + expect(chat.toggleReaction).not.toHaveBeenCalled(); + expect(socket.send).not.toHaveBeenCalled(); + }); + + it.each([ + ["an unknown lobby type", { type: "global" }], + ["a room id that is not a string", { id: 1 }], + ["a missing message id", { messageId: undefined }], + ])("ignores %s", async (_, overrides) => { + const socket = client(); + + await gateway.react(reaction(overrides) as any, socket); + + expect(chat.toggleReaction).not.toHaveBeenCalled(); + expect(socket.send).not.toHaveBeenCalled(); + }); + + it("ignores a missing payload", async () => { + await gateway.react(undefined as any, client()); + + expect(chat.toggleReaction).not.toHaveBeenCalled(); + }); + + it.each([ + ["not on the list", "party"], + ["not a string", 5], + ["missing", undefined], + ])( + "answers a reaction that is %s with invalid, without asking the service", + async (_, value) => { + const socket = client(); + + await gateway.react( + reaction({ reaction: value, requestId: "r-2" }) as any, + socket, + ); + + expect(chat.toggleReaction).not.toHaveBeenCalled(); + expect(sent(socket)).toEqual([ + { + event: "chat:error", + data: { + code: ChatErrorCode.Invalid, + action: "react", + requestId: "r-2", + }, + }, + ]); + }, + ); + + it("asks the service to toggle as the signed in player", async () => { + await gateway.react(reaction({ reaction: "laugh" }) as any, client()); + + expect(chat.toggleReaction).toHaveBeenCalledWith( + ChatLobbyType.Match, + "m-1", + MESSAGE_ID, + "laugh", + expect.objectContaining({ steam_id: "1" }), + ); + }); + + it("acks a toggle under the requestId it came with", async () => { + const socket = client(); + + await gateway.react(reaction({ requestId: "r-3" }) as any, socket); + + expect(sent(socket)).toEqual([ + { + event: "chat:ack", + data: { requestId: "r-3", messageId: MESSAGE_ID, action: "react" }, + }, + ]); + }); + + it("stays quiet for a toggle without a requestId", async () => { + const socket = client(); + + await gateway.react(reaction() as any, socket); + + expect(socket.send).not.toHaveBeenCalled(); + }); + + it.each([ + ChatErrorCode.RateLimited, + ChatErrorCode.NotAllowed, + ChatErrorCode.NotFound, + ChatErrorCode.Gagged, + ])("reports %s under the requestId it came with", async (code) => { + chat.toggleReaction.mockResolvedValue({ toggled: false, code }); + const socket = client(); + + await gateway.react(reaction({ requestId: "r-4" }) as any, socket); + + expect(sent(socket)).toEqual([ + { + event: "chat:error", + data: { code, action: "react", requestId: "r-4" }, + }, + ]); + }); + + it("never relays a reaction to the game server", async () => { + await gateway.react(reaction() as any, client()); + + expect(chat.toggleReaction).toHaveBeenCalled(); + expect(chat.sendChatToServer).not.toHaveBeenCalled(); + }); +}); diff --git a/src/chat/chat.gateway.ts b/src/chat/chat.gateway.ts index 1a162e04..8b9fc01e 100644 --- a/src/chat/chat.gateway.ts +++ b/src/chat/chat.gateway.ts @@ -252,6 +252,56 @@ export class ChatGateway { } } + @SubscribeMessage("lobby:react") + async react( + @MessageBody() + data: { + id: string; + type: ChatLobbyType; + messageId: string; + reaction: unknown; + requestId?: string; + }, + @ConnectedSocket() client: FiveStackWebSocketClient, + ) { + if (!client.user) { + return; + } + + if ( + !ChatGateway.isLobbyType(data?.type) || + typeof data.id !== "string" || + typeof data.messageId !== "string" + ) { + return; + } + + const requestId = + typeof data.requestId === "string" ? data.requestId : undefined; + + if (!ChatService.isReaction(data.reaction)) { + this.sendError(client, "react", ChatErrorCode.Invalid, requestId); + return; + } + + const result = await this.chat.toggleReaction( + data.type, + data.id, + data.messageId, + data.reaction, + client.user, + ); + + if (result.toggled === false) { + this.sendError(client, "react", result.code, requestId); + return; + } + + if (requestId) { + this.sendAck(client, "react", requestId, data.messageId); + } + } + private static isLobbyType(value: unknown): value is ChatLobbyType { return Object.values(ChatLobbyType).includes(value as ChatLobbyType); } diff --git a/src/chat/chat.service.spec.ts b/src/chat/chat.service.spec.ts index 9c5b1bd4..ab2e8f60 100644 --- a/src/chat/chat.service.spec.ts +++ b/src/chat/chat.service.spec.ts @@ -14,10 +14,11 @@ describe("ChatService direct messages", () => { hset: jest.fn(), hget: jest.fn().mockResolvedValue(null), hgetall: jest.fn().mockResolvedValue({}), - hdel: jest.fn(), + hdel: jest.fn().mockResolvedValue(1), get: jest.fn().mockResolvedValue(null), set: jest.fn(), del: jest.fn(), + keys: jest.fn().mockResolvedValue([]), expire: jest.fn(), zadd: jest.fn(), zrevrange: jest.fn().mockResolvedValue([]), @@ -39,6 +40,8 @@ describe("ChatService direct messages", () => { let gagged: boolean; let audited: boolean; let editAuditIds: string[]; + let directReactions: Record | null; + let directReactionFailure: Error | undefined; // The one direct message the fake database holds, if a test put one there. let directMessage: | { @@ -70,6 +73,26 @@ describe("ChatService direct messages", () => { return [{ deleted: audited }]; } + if (sql.includes("FOR NO KEY UPDATE")) { + return directMessage?.id === bindings[0] && + directMessage.roomId === bindings[1] + ? [{ locked: 1 }] + : []; + } + + if (sql.includes("INSERT INTO public.direct_message_reactions")) { + if (directReactionFailure) { + throw directReactionFailure; + } + return []; + } + + if (sql.includes("SELECT reactions.reactions")) { + return directMessage?.id === bindings[0] + ? [{ reactions: directReactions }] + : []; + } + if (sql.includes("INSERT INTO public.chat_message_edits")) { const id = `edit-audit-${editAuditIds.length + 1}`; editAuditIds.push(id); @@ -127,6 +150,17 @@ describe("ChatService direct messages", () => { }), }; + // Every statement in a transaction goes through the same fake as the rest. + Object.assign(postgres, { + transaction: jest.fn(async (work: (client: any) => Promise) => + work({ + query: async (sql: string, bindings: any[]) => ({ + rows: await postgres.query(sql, bindings), + }), + }), + ), + }); + const notifications = { notifyPlayers: jest.fn(), collapseOlderUnread: jest.fn(), @@ -328,6 +362,7 @@ describe("ChatService direct messages", () => { // room would otherwise leave them seated for every test after it. redis.hget.mockResolvedValue(null); redis.hgetall.mockResolvedValue({}); + redis.hdel.mockResolvedValue(1); redis.get.mockResolvedValue(null); redis.eval.mockResolvedValue([1, 1]); notifications.retractChatMessage.mockResolvedValue(undefined); @@ -348,6 +383,8 @@ describe("ChatService direct messages", () => { gagged = false; audited = false; editAuditIds = []; + directReactions = null; + directReactionFailure = undefined; rcon.send.mockResolvedValue(undefined); rcon.connect.mockResolvedValue(rcon); @@ -2195,6 +2232,618 @@ describe("ChatService direct messages", () => { }); }); + describe("reactions", () => { + const MESSAGE_ID = "7b1c2d3e-4f5a-4b6c-8d7e-9f0a1b2c3d4e"; + const ROOM = "chat_match_m-1"; + const REACTIONS = "chat_reactions_match_m-1"; + + let hashes: Record>; + let rates: Record; + let receipts: Set; + + const as = (steamId: string, overrides: Record = {}) => + ({ + steam_id: steamId, + name: "Someone", + role: "user", + ...overrides, + }) as any; + + const react = ( + reaction: unknown = "heart", + user = as(ME), + type = ChatLobbyType.Match, + id = "m-1", + messageId = MESSAGE_ID, + ) => service.toggleReaction(type, id, messageId, reaction, user); + + const broadcasts = (event: string) => + redis.publish.mock.calls + .map(([, payload]) => JSON.parse(payload)) + .filter((published) => published.event === event); + + const toggles = () => + redis.eval.mock.calls.filter(([script]) => script.includes("cjson")); + + const flush = () => new Promise((resolve) => setImmediate(resolve)); + + beforeEach(() => { + rates = {}; + receipts = new Set(); + hashes = { + [ROOM]: { + [MESSAGE_ID]: JSON.stringify({ + id: MESSAGE_ID, + message: "gg", + timestamp: new Date().toISOString(), + source: "web", + from: { role: "user", name: "Friend", steam_id: FRIEND }, + }), + }, + }; + + // Everyone asked about is seated in the room; who may be in it at all + // is the joining tests' subject. + redis.hget.mockImplementation(async (key: string, field: string) => { + if (key.startsWith("chat:")) { + return JSON.stringify({ user: { steam_id: field } }); + } + return hashes[key]?.[field] ?? null; + }); + redis.hgetall.mockImplementation(async (key: string) => { + if (key === "chat:match:m-1") { + return { [FRIEND]: JSON.stringify({ user: { steam_id: FRIEND } }) }; + } + return { ...hashes[key] }; + }); + redis.hdel.mockImplementation(async (key: string, field: string) => { + delete hashes[key]?.[field]; + return 1; + }); + redis.eval.mockImplementation( + async (script: string, _keys: number, ...args: any[]) => { + if (script.includes("INCR")) { + const [key] = args; + rates[key] = (rates[key] ?? 0) + 1; + return rates[key]; + } + + if (script.includes("cjson")) { + const [ + roomKey, + reactionsKey, + receipt, + messageId, + reaction, + steamId, + removeOnly, + ] = args; + + if (receipts.has(receipt)) { + return hashes[reactionsKey]?.[messageId] ?? "{}"; + } + + if (!hashes[roomKey]?.[messageId]) { + return null; + } + + const state = JSON.parse(hashes[reactionsKey]?.[messageId] ?? "{}"); + const holders: string[] = state[reaction] ?? []; + + if (!holders.includes(steamId) && removeOnly === "1") { + return 0; + } + + receipts.add(receipt); + state[reaction] = holders.includes(steamId) + ? holders.filter((holder) => holder !== steamId) + : [...holders, steamId]; + if (state[reaction].length === 0) { + delete state[reaction]; + } + + hashes[reactionsKey] = { + ...hashes[reactionsKey], + [messageId]: JSON.stringify(state), + }; + return JSON.stringify(state); + } + + return [1, 1]; + }, + ); + }); + + it("toggles a reaction on and then off", async () => { + await expect(react()).resolves.toEqual({ + toggled: true, + reactions: { heart: [ME] }, + }); + await expect(react()).resolves.toEqual({ + toggled: true, + reactions: {}, + }); + }); + + it("sends the room the message's whole reaction state", async () => { + await react("heart", as(FRIEND)); + await react("fire", as(FRIEND)); + await react("heart"); + await flush(); + + expect(broadcasts("lobby:match:m-1:reaction").at(-1)).toEqual({ + steamId: FRIEND, + event: "lobby:match:m-1:reaction", + data: { + id: MESSAGE_ID, + reactions: { heart: [FRIEND, ME], fire: [FRIEND] }, + }, + }); + }); + + it("hands the toggle both hashes, a receipt of its own and who reacted", async () => { + await react("laugh"); + await react("laugh"); + + const [first, second] = toggles(); + + expect(first.slice(1)).toEqual([ + 3, + ROOM, + REACTIONS, + expect.stringMatching(/^chat_reaction_applied:[0-9a-f-]{36}$/), + MESSAGE_ID, + "laugh", + ME, + "0", + 300_000, + ]); + expect(second[4]).not.toBe(first[4]); + }); + + it("orders reactions as the list does, whatever order they were stored in", async () => { + await react("sad"); + await react("fire", as(FRIEND)); + + const result = await react("thumbsup"); + + expect(Object.keys(result.toggled ? result.reactions : {})).toEqual([ + "thumbsup", + "fire", + "sad", + ]); + }); + + it("accepts every reaction on the list", async () => { + for (const reaction of ChatService.REACTIONS) { + await expect(react(reaction)).resolves.toMatchObject({ + toggled: true, + }); + } + + expect(ChatService.REACTIONS).toEqual([ + "thumbsup", + "heart", + "laugh", + "fire", + "wow", + "sad", + ]); + }); + + it.each([ + ["one that is not on the list", "party"], + ["a different case", "HEART"], + ["an object key", "__proto__"], + ["an empty string", ""], + ["a number", 5], + ["null", null], + ])("refuses %s as invalid before anything else", async (_, reaction) => { + await expect(react(reaction)).resolves.toEqual({ + toggled: false, + code: ChatErrorCode.Invalid, + }); + + expect(redis.eval).not.toHaveBeenCalled(); + expect(hashes[REACTIONS]).toBeUndefined(); + }); + + it("answers not_found for an id that could never be a message", async () => { + await expect( + react("heart", as(ME), ChatLobbyType.Match, "m-1", "not-a-uuid"), + ).resolves.toEqual({ toggled: false, code: ChatErrorCode.NotFound }); + expect(toggles()).toHaveLength(0); + }); + + it("gives a message that is gone no reactions", async () => { + delete hashes[ROOM][MESSAGE_ID]; + + await expect(react()).resolves.toEqual({ + toggled: false, + code: ChatErrorCode.NotFound, + }); + await flush(); + + expect(hashes[REACTIONS]).toBeUndefined(); + expect(broadcasts("lobby:match:m-1:reaction")).toEqual([]); + }); + + it("refuses someone who cannot get into the room", async () => { + hashes["chat_match_m-2"] = hashes[ROOM]; + + await expect( + react("heart", as(ME), ChatLobbyType.Match, "m-2"), + ).resolves.toEqual({ toggled: false, code: ChatErrorCode.NotAllowed }); + expect(toggles()).toHaveLength(0); + }); + + it("refuses someone who is not in the room", async () => { + redis.hget.mockImplementation(async (key: string, field: string) => + key.startsWith("chat:") ? null : (hashes[key]?.[field] ?? null), + ); + + await expect(react()).resolves.toEqual({ + toggled: false, + code: ChatErrorCode.NotAllowed, + }); + expect(toggles()).toHaveLength(0); + }); + + it("keeps a gagged player from adding a reaction in a group room", async () => { + gagged = true; + + await expect(react()).resolves.toEqual({ + toggled: false, + code: ChatErrorCode.Gagged, + }); + await flush(); + + expect(toggles()[0][8]).toBe("1"); + expect(hashes[REACTIONS]).toBeUndefined(); + expect(broadcasts("lobby:match:m-1:reaction")).toEqual([]); + }); + + it("lets a gagged player take back a reaction they already gave", async () => { + await react("heart"); + gagged = true; + + await expect(react("heart")).resolves.toEqual({ + toggled: true, + reactions: {}, + }); + }); + + it("lets through eight toggles a second and refuses the ninth", async () => { + for (let toggle = 0; toggle < ChatService.REACTION_RATE_LIMIT; toggle++) { + await expect(react()).resolves.toMatchObject({ toggled: true }); + } + + await expect(react()).resolves.toEqual({ + toggled: false, + code: ChatErrorCode.RateLimited, + }); + await flush(); + + expect(toggles()).toHaveLength(ChatService.REACTION_RATE_LIMIT); + expect(broadcasts("lobby:match:m-1:reaction")).toHaveLength( + ChatService.REACTION_RATE_LIMIT, + ); + await expect(react("heart", as(FRIEND))).resolves.toMatchObject({ + toggled: true, + }); + }); + + it("counts toggles per player, in a window that starts with the first", async () => { + await react(); + + const rate = redis.eval.mock.calls.find(([script]) => + script.includes("INCR"), + ); + + expect(rate.slice(1)).toEqual([1, `chat:reaction-rate:${ME}`, 1_000]); + expect(rate[0]).toContain("PEXPIRE"); + }); + + it("never notifies anyone or relays to the game", async () => { + await react(); + await flush(); + + expect(notifications.notifyPlayers).not.toHaveBeenCalled(); + expect(rcon.connect).not.toHaveBeenCalled(); + expect(broadcasts("lobby:match:m-1:chat")).toEqual([]); + }); + + it("clears a message's reactions when it is deleted", async () => { + await react(); + role = "moderator"; + + await expect( + service.deleteMessage( + ChatLobbyType.Match, + "m-1", + MESSAGE_ID, + as(ME, { role: "moderator" }), + ), + ).resolves.toEqual({ deleted: true }); + + expect(redis.hdel).toHaveBeenCalledWith(REACTIONS, MESSAGE_ID); + expect(hashes[REACTIONS][MESSAGE_ID]).toBeUndefined(); + }); + + it("clears a message's reactions when its author deletes it", async () => { + hashes[ROOM][MESSAGE_ID] = JSON.stringify({ + ...JSON.parse(hashes[ROOM][MESSAGE_ID]), + from: { role: "user", name: "Me", steam_id: ME }, + }); + await react("heart", as(FRIEND)); + + await expect( + service.deleteMessage(ChatLobbyType.Match, "m-1", MESSAGE_ID, as(ME)), + ).resolves.toEqual({ deleted: true }); + + expect(hashes[REACTIONS][MESSAGE_ID]).toBeUndefined(); + }); + + it("still announces a delete when its reactions cannot be cleared", async () => { + await react(); + role = "moderator"; + redis.hdel.mockImplementation(async (key: string, field: string) => { + if (key === REACTIONS) { + throw new Error("connection reset"); + } + delete hashes[key]?.[field]; + return 1; + }); + + await expect( + service.deleteMessage( + ChatLobbyType.Match, + "m-1", + MESSAGE_ID, + as(ME, { role: "moderator" }), + ), + ).resolves.toEqual({ deleted: true }); + await flush(); + + expect(broadcasts("lobby:match:m-1:deleted")).toHaveLength(1); + expect(notifications.retractChatMessage).toHaveBeenCalledWith(MESSAGE_ID); + }); + + it("puts each message's reactions in the room's history", async () => { + const quiet = "8c2d3e4f-5a6b-4c7d-9e8f-0a1b2c3d4e5f"; + hashes[ROOM][quiet] = JSON.stringify({ + id: quiet, + message: "later", + timestamp: new Date(Date.now() + 1_000).toISOString(), + source: "web", + from: { role: "user", name: "Friend", steam_id: FRIEND }, + }); + await react("fire"); + + const history = await service["getMessages"](ChatLobbyType.Match, "m-1"); + + expect( + history.map(({ id, reactions }: any) => ({ id, reactions })), + ).toEqual([ + { id: MESSAGE_ID, reactions: { fire: [ME] } }, + { id: quiet, reactions: {} }, + ]); + }); + + it("sends a new message out with no reactions, and stores none", async () => { + await service.sendMessageToChat(ChatLobbyType.Match, "m-1", as(ME), "hi"); + await flush(); + + const [chat] = broadcasts("lobby:match:m-1:chat"); + const stored = JSON.parse( + redis.hset.mock.calls.find(([key]) => key === ROOM)[2], + ); + + expect(chat.data.reactions).toEqual({}); + expect(stored).not.toHaveProperty("reactions"); + }); + + // The move itself is a Lua script, exercised against real redis in + // test/chat-redis-actions.spec.ts; this is what the service does around it. + it("hands the move all four keys, then re-sends history with the reactions it carried", async () => { + hashes[REACTIONS] = { [MESSAGE_ID]: JSON.stringify({ fire: [ME] }) }; + redis.eval.mockResolvedValueOnce(1); + + await service.migrateLobbyMessages( + ChatLobbyType.Draft, + "d-1", + ChatLobbyType.Match, + "m-1", + ); + await flush(); + + const move = redis.eval.mock.calls.find(([script]) => + script.includes("HGETALL"), + ); + + expect(move.slice(1)).toEqual([ + 4, + "chat_draft_d-1", + ROOM, + "chat_reactions_draft_d-1", + REACTIONS, + 3600, + ]); + expect(broadcasts("lobby:match:m-1:messages")).toEqual([ + expect.objectContaining({ + data: { + id: "m-1", + messages: [ + expect.objectContaining({ + id: MESSAGE_ID, + reactions: { fire: [ME] }, + }), + ], + }, + }), + ]); + }); + + describe("in a direct conversation", () => { + const room = directRoomId(ME, FRIEND); + + const reactDirect = (reaction = "heart", user = as(ME)) => + react(reaction, user, ChatLobbyType.Direct, room); + + beforeEach(() => { + directMessage = { + id: MESSAGE_ID, + roomId: room, + author: FRIEND, + message: "gg", + open: true, + editedAt: null, + }; + directReactions = { heart: [ME] }; + redis.hgetall.mockImplementation(async (key: string) => + key === `chat:direct:${room}` + ? { [FRIEND]: JSON.stringify({ user: { steam_id: FRIEND } }) } + : {}, + ); + }); + + it("locks the message, toggles in one statement scoped to the conversation, then reads the state", async () => { + await expect(reactDirect()).resolves.toEqual({ + toggled: true, + reactions: { heart: [ME] }, + }); + await flush(); + + const lock = queries.findIndex(({ sql }) => + sql.includes("FOR NO KEY UPDATE"), + ); + const toggleAt = queries.findIndex(({ sql }) => + sql.includes("INSERT INTO public.direct_message_reactions"), + ); + const toggle = queries[toggleAt]; + + expect(postgres.transaction).toHaveBeenCalledTimes(1); + expect(queries[lock].bindings).toEqual([MESSAGE_ID, room]); + expect(lock).toBeLessThan(toggleAt); + expect(toggle.bindings).toEqual([MESSAGE_ID, ME, "heart", room]); + expect(toggle.sql).toContain( + "DELETE FROM public.direct_message_reactions", + ); + expect(broadcasts(`lobby:direct:${room}:reaction`)).toEqual([ + { + steamId: FRIEND, + event: `lobby:direct:${room}:reaction`, + data: { id: MESSAGE_ID, reactions: { heart: [ME] } }, + }, + ]); + }); + + it("answers with an empty state once the last reaction is gone", async () => { + directReactions = null; + + await expect(reactDirect()).resolves.toEqual({ + toggled: true, + reactions: {}, + }); + }); + + it("is not held back by a gag", async () => { + gagged = true; + + await expect(reactDirect()).resolves.toMatchObject({ toggled: true }); + expect( + queries.some(({ sql }) => sql.includes("public.is_gagged")), + ).toBe(false); + }); + + it("refuses once the friendship is gone", async () => { + acceptedFriendships = []; + + await expect(reactDirect()).resolves.toEqual({ + toggled: false, + code: ChatErrorCode.NotAllowed, + }); + expect( + queries.some(({ sql }) => sql.includes("direct_message_reactions")), + ).toBe(false); + }); + + it("answers not_found for a message that is not in the conversation", async () => { + directMessage = { + ...directMessage, + roomId: directRoomId(FRIEND, STRANGER), + }; + + await expect(reactDirect()).resolves.toEqual({ + toggled: false, + code: ChatErrorCode.NotFound, + }); + await flush(); + + expect( + queries.some(({ sql }) => + sql.includes("INSERT INTO public.direct_message_reactions"), + ), + ).toBe(false); + expect(broadcasts(`lobby:direct:${room}:reaction`)).toEqual([]); + }); + + it("orders reactions as the list does", async () => { + directReactions = { sad: [FRIEND], heart: [ME], thumbsup: [FRIEND] }; + + const result = await reactDirect(); + + expect(Object.keys(result.toggled ? result.reactions : {})).toEqual([ + "thumbsup", + "heart", + "sad", + ]); + }); + + it("lets any other failure through", async () => { + directReactionFailure = new Error("database down"); + + await expect(reactDirect()).rejects.toThrow("database down"); + }); + + it("puts reactions in the conversation's history", async () => { + postgres.query.mockResolvedValueOnce([ + { + id: "dm-1", + message: "old", + created_at: new Date("2026-01-01T00:00:00Z"), + reactions: { heart: [FRIEND] }, + steam_id: ME, + name: "Someone", + role: "user", + avatar_url: null, + profile_url: null, + }, + { + id: "dm-2", + message: "older", + created_at: new Date("2025-12-31T00:00:00Z"), + reactions: null, + steam_id: ME, + name: "Someone", + role: "user", + avatar_url: null, + profile_url: null, + }, + ]); + + const history = await service["getDirectMessages"](room); + + expect( + history.map(({ id, reactions }: any) => ({ id, reactions })), + ).toEqual([ + { id: "dm-2", reactions: {} }, + { id: "dm-1", reactions: { heart: [FRIEND] } }, + ]); + }); + }); + }); + describe("rosters", () => { it("resolves both parties of a conversation", async () => { expect( diff --git a/src/chat/chat.service.ts b/src/chat/chat.service.ts index 840c0d9c..ca22231f 100644 --- a/src/chat/chat.service.ts +++ b/src/chat/chat.service.ts @@ -23,6 +23,8 @@ import { ChatMessage, ChatMessageSource } from "./types/ChatMessage"; import { ChatSendResult } from "./types/ChatSendResult"; import { ChatDeleteResult } from "./types/ChatDeleteResult"; import { ChatEditResult } from "./types/ChatEditResult"; +import { ChatReactions } from "./types/ChatReactions"; +import { ChatReactResult } from "./types/ChatReactResult"; @Injectable() export class ChatService { @@ -90,19 +92,127 @@ export class ChatService { private static readonly EDIT_RECEIPT_TTL_MS = 60 * 60 * 1000; - // One step, so an edit or a delete in the old room lands either before the - // move and is carried with it, or after it and finds nothing -- never on a - // copy that is about to be thrown away or written back. + // One step, so an edit, a delete or a reaction in the old room lands either + // before the move and is carried with it, or after it and finds nothing -- + // never on a copy that is about to be thrown away or written back. + // + // A room TTL of 0 drops a moved message at once, and its reactions must not + // then be left behind with no expiry at all. private static readonly MOVE_ROOM_MESSAGES_SCRIPT = ` local messages = redis.call('HGETALL', KEYS[1]) for i = 1, #messages, 2 do redis.call('HSET', KEYS[2], messages[i], messages[i + 1]) redis.call('HEXPIRE', KEYS[2], ARGV[1], 'FIELDS', 1, messages[i]) + local reactions = redis.call('HGET', KEYS[3], messages[i]) + local expiresAt = redis.call('HPEXPIRETIME', KEYS[2], 'FIELDS', 1, messages[i])[1] + if reactions and expiresAt ~= -2 then + redis.call('HSET', KEYS[4], messages[i], reactions) + if expiresAt > 0 then + redis.call('HPEXPIREAT', KEYS[4], expiresAt, 'FIELDS', 1, messages[i]) + end + end end - redis.call('DEL', KEYS[1]) + redis.call('DEL', KEYS[1], KEYS[3]) return #messages / 2 `; + // The web maps each id to its glyph; laugh is 😂. + public static readonly REACTIONS = [ + "thumbsup", + "heart", + "laugh", + "fire", + "wow", + "sad", + ] as const; + + // Every toggle is a broadcast to the whole room, so this is what keeps one + // client from flooding everyone else in it. + public static readonly REACTION_RATE_LIMIT = 8; + + private static readonly REACTION_RATE_WINDOW_MS = 1_000; + + private static readonly REACTION_RATE_SCRIPT = ` + local count = redis.call('INCR', KEYS[1]) + if count == 1 then + redis.call('PEXPIRE', KEYS[1], ARGV[1]) + end + return count + `; + + // Reactions live beside the messages rather than inside their JSON, so they + // never contend with an edit's compare-and-set. A message that is gone gets + // none: they would outlive it with no expiry. HSET drops the field's expiry, + // so the reactions are given the message's own, and go when it does. + // + // KEYS[3] is a receipt for this one toggle: ioredis resends a command whose + // reply was lost to a reconnect, and a toggle run twice undoes itself. + // ARGV[4] = '1' allows only taking a reaction back (a gagged player), and + // 0 is the answer when this toggle would have added one. + private static readonly TOGGLE_ROOM_REACTION_SCRIPT = ` + if redis.call('EXISTS', KEYS[3]) == 1 then + return redis.call('HGET', KEYS[2], ARGV[1]) or '{}' + end + if redis.call('HEXISTS', KEYS[1], ARGV[1]) == 0 then + return false + end + local stored = redis.call('HGET', KEYS[2], ARGV[1]) + local reactions = {} + if stored then + reactions = cjson.decode(stored) + end + local reactors = {} + local removed = false + for _, steamId in ipairs(reactions[ARGV[2]] or {}) do + if steamId == ARGV[3] then + removed = true + else + table.insert(reactors, steamId) + end + end + if not removed then + if ARGV[4] == '1' then + return 0 + end + table.insert(reactors, ARGV[3]) + end + if #reactors == 0 then + reactions[ARGV[2]] = nil + else + reactions[ARGV[2]] = reactors + end + redis.call('SET', KEYS[3], '1', 'PX', ARGV[5]) + if next(reactions) == nil then + redis.call('HDEL', KEYS[2], ARGV[1]) + return '{}' + end + local encoded = cjson.encode(reactions) + redis.call('HSET', KEYS[2], ARGV[1], encoded) + local expiresAt = redis.call('HPEXPIRETIME', KEYS[1], 'FIELDS', 1, ARGV[1])[1] + if expiresAt > 0 then + redis.call('HPEXPIREAT', KEYS[2], expiresAt, 'FIELDS', 1, ARGV[1]) + end + return encoded + `; + + private static readonly REACTION_RECEIPT_TTL_MS = 5 * 60 * 1000; + + // Each direct message's reactions as one object, for a query that has the + // message as `dm`. + private static readonly DIRECT_MESSAGE_REACTIONS = `LEFT JOIN LATERAL ( + SELECT jsonb_object_agg(grouped.reaction, grouped.steam_ids) + AS reactions + FROM ( + SELECT r.reaction, + jsonb_agg(r.steam_id::text + ORDER BY r.created_at, r.steam_id) + AS steam_ids + FROM public.direct_message_reactions r + WHERE r.message_id = dm.id + GROUP BY r.reaction + ) grouped + ) reactions ON true`; + // Shared by every direct message edit and delete, so the author and window // are judged in the same statement that changes the row, on the database's // clock -- the one created_at was stamped with. @@ -483,16 +593,52 @@ export class ChatService { return await this.getDirectMessages(id); } - const messagesObject = await this.redis.hgetall(`chat_${type}_${id}`); + return await this.getRoomMessages(type, id); + } - return Object.values(messagesObject) - .map((value) => JSON.parse(value)) + // Both hashes are asked for at once, so this is still one round trip. A + // reaction landing between the two reads is harmless: the room is sent the + // message's whole reaction state every time one changes. + private async getRoomMessages( + type: ChatLobbyType, + id: string, + ): Promise { + const [messages, reactions] = await Promise.all([ + this.redis.hgetall(`chat_${type}_${id}`), + this.redis.hgetall(ChatService.reactionsKey(type, id)), + ]); + + return Object.entries(messages) + .map( + ([messageId, value]): ChatMessage => ({ + ...(JSON.parse(value) as ChatMessage), + reactions: ChatService.orderedReactions( + reactions[messageId] ? JSON.parse(reactions[messageId]) : null, + ), + }), + ) .sort( (a, b) => new Date(a.timestamp).getTime() - new Date(b.timestamp).getTime(), ); } + // The stores keep reactions in whatever order they like (a Lua table's, + // jsonb's by key length), so clients are always handed the list's order. + private static orderedReactions( + state: ChatReactions | null | undefined, + ): ChatReactions { + return Object.fromEntries( + ChatService.REACTIONS.filter( + (reaction) => (state?.[reaction]?.length ?? 0) > 0, + ).map((reaction) => [reaction, state[reaction]]), + ); + } + + private static reactionsKey(type: ChatLobbyType, id: string) { + return `chat_reactions_${type}_${id}`; + } + private async refreshClientUser(client: FiveStackWebSocketClient) { if (!client.user?.steam_id) { return; @@ -690,10 +836,12 @@ export class ChatService { ); } - void this.to(type, id, "chat", message); + const outgoing: ChatMessage = { ...message, reactions: {} }; + + void this.to(type, id, "chat", outgoing); if (type === ChatLobbyType.Direct) { - void this.deliverDirectMessage(id, player, message); + void this.deliverDirectMessage(id, player, outgoing); } // Best effort, and never allowed to take a message delivery down with it. @@ -790,6 +938,18 @@ export class ChatService { await this.redis.hdel(messageKey, messageId); + // Never allowed to stop the delete being announced: left behind, the + // reactions still expire with the message, and nothing reads reactions + // for a message that is gone. + await this.redis + .hdel(ChatService.reactionsKey(type, id), messageId) + .catch((error) => { + this.logger.warn( + `unable to clear reactions for ${type}:${id} message ${messageId}`, + error, + ); + }); + void this.to(type, id, "deleted", { id: messageId }); await this.retractNotifications(type, id, messageId); @@ -882,6 +1042,165 @@ export class ChatService { ); } + public static isReaction( + value: unknown, + ): value is (typeof ChatService.REACTIONS)[number] { + return (ChatService.REACTIONS as readonly unknown[]).includes(value); + } + + // Never a notification and never relayed to the game: the room is sent the + // message's whole reaction state, so a client that missed a toggle is right + // again with the next one. + public async toggleReaction( + type: ChatLobbyType, + id: string, + messageId: string, + reaction: unknown, + user: User, + ): Promise { + if (!ChatService.isReaction(reaction)) { + return { toggled: false, code: ChatErrorCode.Invalid }; + } + + if (!ChatService.UUID.test(messageId)) { + return { toggled: false, code: ChatErrorCode.NotFound }; + } + + if (!(await this.withinReactionRate(user.steam_id))) { + return { toggled: false, code: ChatErrorCode.RateLimited }; + } + + if (!(await this.canPostIn(type, id, user))) { + return { toggled: false, code: ChatErrorCode.NotAllowed }; + } + + const reactions = + type === ChatLobbyType.Direct + ? await this.toggleDirectReaction(id, messageId, reaction, user) + : await this.toggleRoomReaction( + type, + id, + messageId, + reaction, + user, + // Like deleting their own message, taking a reaction back is not + // speech, so a gag leaves it alone. + await this.isGagged(user.steam_id), + ); + + if (reactions === "gagged") { + return { toggled: false, code: ChatErrorCode.Gagged }; + } + + if (!reactions) { + return { toggled: false, code: ChatErrorCode.NotFound }; + } + + void this.to(type, id, "reaction", { id: messageId, reactions }); + + return { toggled: true, reactions }; + } + + private async withinReactionRate(steamId: string): Promise { + const count = await this.redis.eval( + ChatService.REACTION_RATE_SCRIPT, + 1, + `chat:reaction-rate:${steamId}`, + ChatService.REACTION_RATE_WINDOW_MS, + ); + + return Number(count) <= ChatService.REACTION_RATE_LIMIT; + } + + private async toggleRoomReaction( + type: ChatLobbyType, + id: string, + messageId: string, + reaction: string, + user: User, + removeOnly: boolean, + ): Promise { + const state = await this.redis.eval( + ChatService.TOGGLE_ROOM_REACTION_SCRIPT, + 3, + `chat_${type}_${id}`, + ChatService.reactionsKey(type, id), + `chat_reaction_applied:${randomUUID()}`, + messageId, + reaction, + String(user.steam_id), + removeOnly ? "1" : "0", + ChatService.REACTION_RECEIPT_TTL_MS, + ); + + if (state === 0) { + return "gagged"; + } + + if (typeof state !== "string") { + return null; + } + + return ChatService.orderedReactions(JSON.parse(state)); + } + + private async toggleDirectReaction( + roomId: string, + messageId: string, + reaction: string, + user: User, + ): Promise { + return await this.postgres.transaction(async (client) => { + // Toggles on one message take turns here, and each then starts its + // statement after the last one committed. Two toggles on an absent row + // in one snapshot would both insert, and the second would quietly do + // nothing instead of taking it back. The lock also holds off the + // message's deletion until the reaction is written. + const { rows: found } = await client.query( + `SELECT 1 FROM public.direct_messages + WHERE id = $1::uuid AND room_id = $2 + FOR NO KEY UPDATE`, + [messageId, roomId], + ); + + if (found.length === 0) { + return null; + } + + await client.query( + `WITH removed AS ( + DELETE FROM public.direct_message_reactions + WHERE message_id = $1::uuid + AND steam_id = $2::bigint + AND reaction = $3 + RETURNING 1 + ) + INSERT INTO public.direct_message_reactions + (message_id, steam_id, reaction) + SELECT $1::uuid, $2::bigint, $3 + WHERE NOT EXISTS (SELECT 1 FROM removed) + AND EXISTS ( + SELECT 1 FROM public.direct_messages + WHERE id = $1::uuid AND room_id = $4 + ) + ON CONFLICT DO NOTHING`, + [messageId, String(user.steam_id), reaction, roomId], + ); + + const { + rows: [row], + } = await client.query<{ reactions: ChatReactions | null }>( + `SELECT reactions.reactions + FROM public.direct_messages dm + ${ChatService.DIRECT_MESSAGE_REACTIONS} + WHERE dm.id = $1::uuid`, + [messageId], + ); + + return ChatService.orderedReactions(row?.reactions); + }); + } + private async editRoomMessage( type: ChatLobbyType, id: string, @@ -1887,6 +2206,7 @@ export class ChatService { message: string; created_at: Date; edited_at: Date | null; + reactions: ChatReactions | null; steam_id: string; name: string; role: e_player_roles_enum; @@ -1895,10 +2215,12 @@ export class ChatService { }> >( `SELECT dm.id::text AS id, dm.message, dm.created_at, dm.edited_at, + reactions.reactions, p.steam_id::text AS steam_id, p.name, p.role::text AS role, p.avatar_url, p.profile_url FROM public.direct_messages dm JOIN public.players p ON p.steam_id = dm.from_steam_id + ${ChatService.DIRECT_MESSAGE_REACTIONS} WHERE dm.room_id = $1 ORDER BY dm.created_at DESC, dm.seq DESC LIMIT 200`, @@ -1915,6 +2237,7 @@ export class ChatService { ...(row.edited_at ? { edited_at: new Date(row.edited_at).toISOString() } : {}), + reactions: ChatService.orderedReactions(row.reactions), from: { role: row.role, name: row.name, @@ -2079,6 +2402,7 @@ export class ChatService { | "chat" | "edited" | "deleted" + | "reaction" | "list" | "messages" | "joined" @@ -2407,13 +2731,13 @@ export class ChatService { toType: ChatLobbyType, toId: string, ) { - const toKey = `chat_${toType}_${toId}`; - const moved = await this.redis.eval( ChatService.MOVE_ROOM_MESSAGES_SCRIPT, - 2, + 4, `chat_${fromType}_${fromId}`, - toKey, + `chat_${toType}_${toId}`, + ChatService.reactionsKey(fromType, fromId), + ChatService.reactionsKey(toType, toId), this.ttlFor(toType), ); @@ -2423,13 +2747,7 @@ export class ChatService { return; } - const merged = await this.redis.hgetall(toKey); - const messages = Object.values(merged) - .map((value) => JSON.parse(value)) - .sort( - (a, b) => - new Date(a.timestamp).getTime() - new Date(b.timestamp).getTime(), - ); + const messages = await this.getRoomMessages(toType, toId); void this.to(toType, toId, "messages", { id: toId, messages }); } diff --git a/src/chat/enums/ChatErrorCode.ts b/src/chat/enums/ChatErrorCode.ts index 00aae33a..38814cf9 100644 --- a/src/chat/enums/ChatErrorCode.ts +++ b/src/chat/enums/ChatErrorCode.ts @@ -7,4 +7,5 @@ export enum ChatErrorCode { Gagged = "gagged", NotFound = "not_found", WindowClosed = "window_closed", + RateLimited = "rate_limited", } diff --git a/src/chat/types/ChatAction.ts b/src/chat/types/ChatAction.ts index 029223aa..0c84bafd 100644 --- a/src/chat/types/ChatAction.ts +++ b/src/chat/types/ChatAction.ts @@ -1,3 +1,3 @@ // Echoed on `chat:ack` and `chat:error` so the client knows which request the // answer is for. A contract with the web -- add to it, never rename. -export type ChatAction = "send" | "delete" | "edit"; +export type ChatAction = "send" | "delete" | "edit" | "react"; diff --git a/src/chat/types/ChatMessage.ts b/src/chat/types/ChatMessage.ts index 0ba56e32..b18b11b1 100644 --- a/src/chat/types/ChatMessage.ts +++ b/src/chat/types/ChatMessage.ts @@ -1,4 +1,5 @@ import { e_player_roles_enum } from "generated"; +import { ChatReactions } from "./ChatReactions"; export type ChatMessageSource = "web" | "game"; @@ -10,6 +11,9 @@ export interface ChatMessage { source?: ChatMessageSource; // ISO 8601, present once the author has edited the message. edited_at?: string; + // Added when a message is sent to clients, never stored in its JSON: an + // edit's compare-and-set would otherwise contend with every reaction. + reactions?: ChatReactions; from: { role: e_player_roles_enum; name: string; diff --git a/src/chat/types/ChatReactResult.ts b/src/chat/types/ChatReactResult.ts new file mode 100644 index 00000000..14fd9bb2 --- /dev/null +++ b/src/chat/types/ChatReactResult.ts @@ -0,0 +1,6 @@ +import { ChatErrorCode } from "../enums/ChatErrorCode"; +import { ChatReactions } from "./ChatReactions"; + +export type ChatReactResult = + | { toggled: true; reactions: ChatReactions } + | { toggled: false; code: ChatErrorCode }; diff --git a/src/chat/types/ChatReactions.ts b/src/chat/types/ChatReactions.ts new file mode 100644 index 00000000..7a2b121f --- /dev/null +++ b/src/chat/types/ChatReactions.ts @@ -0,0 +1,3 @@ +// Reaction id to the steam ids that reacted with it, oldest first. A reaction +// nobody holds is left out rather than sent as an empty list. +export type ChatReactions = Record; diff --git a/test/chat-direct-messages.spec.ts b/test/chat-direct-messages.spec.ts index 39b615ab..f179ba43 100644 --- a/test/chat-direct-messages.spec.ts +++ b/test/chat-direct-messages.spec.ts @@ -653,6 +653,241 @@ describe("direct messages (SQL-driven)", () => { }); }); + describe("reactions", () => { + const as = (steamId: string) => ({ steam_id: steamId }) as any; + + const sent = async (roomId: string, from: string, message = "gg") => { + const result = await say(roomId, from, message); + return result.accepted ? result.messageId : ""; + }; + + const react = ( + roomId: string, + id: string, + steamId: string, + reaction = "heart", + ) => + chat.toggleReaction( + ChatLobbyType.Direct, + roomId, + id, + reaction, + as(steamId), + ); + + const rows = () => + postgres.query>( + `SELECT message_id::text AS message_id FROM direct_message_reactions`, + ); + + // Reacting is held to the same rule as sending, so both have to be seated + // in the room, and the rate limit wants a real count back. + beforeEach(() => { + redis.hget.mockResolvedValue(JSON.stringify({ user: {} })); + redis.eval.mockImplementation(async (script: string) => + script.includes("INCR") ? 1 : [1, 1], + ); + }); + + afterEach(() => { + redis.hget.mockResolvedValue(null); + redis.eval.mockResolvedValue([1, 1]); + }); + + it("toggles each player's reactions, oldest first", async () => { + const me = await fx.player(); + const friend = await fx.player(); + const room = directRoomId(me, friend); + const id = await sent(room, friend); + + await expect(react(room, id, me)).resolves.toEqual({ + toggled: true, + reactions: { heart: [me] }, + }); + await react(room, id, friend); + await expect(react(room, id, me, "laugh")).resolves.toEqual({ + toggled: true, + reactions: { heart: [me, friend], laugh: [me] }, + }); + await expect(react(room, id, me)).resolves.toEqual({ + toggled: true, + reactions: { heart: [friend], laugh: [me] }, + }); + await react(room, id, friend); + await expect(react(room, id, me, "laugh")).resolves.toEqual({ + toggled: true, + reactions: {}, + }); + + expect(await rows()).toEqual([]); + }); + + it("cancels one player's racing toggles out in pairs", async () => { + const me = await fx.player(); + const friend = await fx.player(); + const room = directRoomId(me, friend); + + // One round loses the race without the row lock only some of the time. + for (let round = 0; round < 8; round++) { + await postgres.query("DELETE FROM direct_message_reactions"); + const even = await sent(room, friend, "even"); + const odd = await sent(room, friend, "odd"); + + await Promise.all([ + ...Array.from({ length: 6 }, () => react(room, even, me)), + ...Array.from({ length: 5 }, () => react(room, odd, me)), + ]); + + expect(await rows()).toEqual([{ message_id: odd }]); + } + }); + + it("counts both parties toggling every reaction at once exactly", async () => { + const me = await fx.player(); + const friend = await fx.player(); + const room = directRoomId(me, friend); + const id = await sent(room, friend); + + const results = await Promise.all( + [me, friend].flatMap((steamId) => + ChatService.REACTIONS.map((reaction) => + react(room, id, steamId, reaction), + ), + ), + ); + + expect(results.every((result) => result.toggled)).toBe(true); + + const [message] = await chat["getMessages"](ChatLobbyType.Direct, room); + + expect(Object.keys(message.reactions)).toEqual([ + ...ChatService.REACTIONS, + ]); + for (const reaction of ChatService.REACTIONS) { + expect([...message.reactions[reaction]].sort()).toEqual( + [me, friend].sort(), + ); + } + }); + + it("hands back reactions with the conversation's history", async () => { + const me = await fx.player(); + const friend = await fx.player(); + const room = directRoomId(me, friend); + const reacted = await sent(room, friend, "first"); + const quiet = await sent(room, me, "second"); + + await react(room, reacted, me, "fire"); + + const history = await chat["getMessages"](ChatLobbyType.Direct, room); + + expect(history.map(({ id, reactions }) => ({ id, reactions }))).toEqual([ + { id: reacted, reactions: { fire: [me] } }, + { id: quiet, reactions: {} }, + ]); + }); + + it("answers not_found for a message from another conversation, and writes nothing", async () => { + const me = await fx.player(); + const friend = await fx.player(); + const other = await fx.player(); + const room = directRoomId(me, friend); + const elsewhere = await sent(directRoomId(friend, other), friend); + + await expect(react(room, elsewhere, me)).resolves.toEqual({ + toggled: false, + code: ChatErrorCode.NotFound, + }); + expect(await rows()).toEqual([]); + }); + + it("goes with the message when its author deletes it", async () => { + const me = await fx.player(); + const friend = await fx.player(); + const room = directRoomId(me, friend); + const id = await sent(room, friend); + const kept = await sent(room, me); + + await react(room, id, me); + await react(room, kept, friend); + await chat.deleteMessage(ChatLobbyType.Direct, room, id, as(friend)); + + expect(await rows()).toEqual([{ message_id: kept }]); + }); + + it("goes with the message when retention sweeps it", async () => { + const me = await fx.player(); + const friend = await fx.player(); + const room = directRoomId(me, friend); + const id = await sent(room, friend); + + await react(room, id, me); + await postgres.query( + `UPDATE direct_messages SET created_at = now() - interval '400 days'`, + ); + + await new PruneDirectMessages(logger as any, postgres).process({} as any); + + expect(await rows()).toEqual([]); + }); + + it("goes with the player who reacted", async () => { + const me = await fx.player(); + const friend = await fx.player(); + const room = directRoomId(me, friend); + const id = await sent(room, friend); + + await react(room, id, me); + await postgres.query(`DELETE FROM players WHERE steam_id = $1::bigint`, [ + me, + ]); + + expect(await rows()).toEqual([]); + }); + + describe("the migration", () => { + const migration = (file: string) => + readFileSync( + join( + __dirname, + "../hasura/migrations/default/1888000000300_direct_message_reactions", + file, + ), + "utf8", + ); + + const table = async () => + ( + await postgres.query>( + `SELECT to_regclass('public.direct_message_reactions')::text AS table`, + ) + )[0].table; + + it("re-applies cleanly and rolls back", async () => { + await postgres.query(migration("up.sql")); + await postgres.query(migration("up.sql")); + expect(await table()).toBe("direct_message_reactions"); + + const indexes = await postgres.query>( + `SELECT indexname FROM pg_indexes + WHERE schemaname = 'public' + AND tablename = 'direct_message_reactions' + ORDER BY indexname`, + ); + expect(indexes.map(({ indexname }) => indexname)).toEqual([ + "direct_message_reactions_pkey", + "direct_message_reactions_steam_id_idx", + ]); + + await postgres.query(migration("down.sql")); + expect(await table()).toBeNull(); + + await postgres.query(migration("up.sql")); + expect(await table()).toBe("direct_message_reactions"); + }); + }); + }); + describe("the edited_at migration", () => { const migration = (file: string) => readFileSync( diff --git a/test/chat-redis-actions.spec.ts b/test/chat-redis-actions.spec.ts index 8ca15907..ad0a5e02 100644 --- a/test/chat-redis-actions.spec.ts +++ b/test/chat-redis-actions.spec.ts @@ -687,6 +687,419 @@ describe("chat edits and self deletes (SQL-driven)", () => { }); }); + describe("reactions", () => { + const reactionsKey = (matchId: string, type = "match") => + `chat_reactions_${type}_${matchId}`; + + const reactionsExpireAt = async ( + matchId: string, + id: string, + type = "match", + ) => + ( + (await redis.call( + "HPEXPIRETIME", + reactionsKey(matchId, type), + "FIELDS", + 1, + id, + )) as number[] + )[0]; + + const player = (index: number) => + ({ + steam_id: String(76561199620000000n + BigInt(index)), + name: `Player ${index}`, + role: "user", + }) as any; + + const seat = (matchId: string, user: any, type = "match") => + redis.hset( + `chat:${type}:${matchId}`, + user.steam_id, + JSON.stringify({ user: { steam_id: user.steam_id } }), + ); + + const react = ( + matchId: string, + id: string, + user: any, + reaction = "heart", + type = ChatLobbyType.Match, + ) => chat.toggleReaction(type, matchId, id, reaction, user); + + const stored = async (matchId: string, id: string, type = "match") => + JSON.parse((await redis.hget(reactionsKey(matchId, type), id)) ?? "null"); + + it("gives reactions the message's own expiry, through an edit", async () => { + const user = await author(); + const matchId = randomUUID(); + await seat(matchId, user); + + const sent = await chat.sendMessageToChat( + ChatLobbyType.Match, + matchId, + user, + "typo", + true, + ); + const id = sent.accepted ? sent.messageId : ""; + const expiry = await expiresAt(matchId, id); + + await react(matchId, id, user); + await new Promise((resolve) => setTimeout(resolve, 20)); + await react(matchId, id, user, "fire"); + + expect(expiry).toBeGreaterThan(Date.now()); + expect(await reactionsExpireAt(matchId, id)).toBe(expiry); + + await edit(matchId, id, user); + + expect(await expiresAt(matchId, id)).toBe(expiry); + expect(await reactionsExpireAt(matchId, id)).toBe(expiry); + expect(await stored(matchId, id)).toEqual({ + heart: [user.steam_id], + fire: [user.steam_id], + }); + }); + + it("keeps an expiry to the millisecond, and none for a message with none", async () => { + const user = await author(); + const matchId = randomUUID(); + await seat(matchId, user); + const timed = await place(matchId, user); + const untimed = await place(matchId, user); + const at = Date.now() + 123_457; + await redis.call("HPEXPIREAT", key(matchId), at, "FIELDS", 1, timed); + + await react(matchId, timed, user); + await react(matchId, untimed, user); + + expect(await reactionsExpireAt(matchId, timed)).toBe(at); + expect(await reactionsExpireAt(matchId, untimed)).toBe(-1); + }); + + it("keeps the message's expiry when a reaction is taken back and others remain", async () => { + const user = await author(); + const other = player(1); + const matchId = randomUUID(); + await seat(matchId, user); + await seat(matchId, other); + const id = await place(matchId, user); + const at = Date.now() + 98_765; + await redis.call("HPEXPIREAT", key(matchId), at, "FIELDS", 1, id); + + await react(matchId, id, user); + await react(matchId, id, other); + await react(matchId, id, other, "fire"); + await react(matchId, id, user); + + expect(await stored(matchId, id)).toEqual({ + heart: [other.steam_id], + fire: [other.steam_id], + }); + expect(await reactionsExpireAt(matchId, id)).toBe(at); + }); + + it("lets reactions go when the message does", async () => { + const user = await author(); + const matchId = randomUUID(); + await seat(matchId, user); + const id = await place(matchId, user); + await redis.call("HPEXPIRE", key(matchId), 150, "FIELDS", 1, id); + + await react(matchId, id, user); + await new Promise((resolve) => setTimeout(resolve, 300)); + + expect(await redis.exists(reactionsKey(matchId))).toBe(0); + }); + + it("counts twenty players reacting at once exactly", async () => { + const matchId = randomUUID(); + const id = await place(matchId, await author()); + const players = Array.from({ length: 20 }, (_, index) => player(index)); + await Promise.all(players.map((user) => seat(matchId, user))); + + const on = await Promise.all( + players.flatMap((user) => [ + react(matchId, id, user, "heart"), + react(matchId, id, user, "fire"), + ]), + ); + + expect(on.every((result) => result.toggled)).toBe(true); + + const state = await stored(matchId, id); + expect([...state.heart].sort()).toEqual( + players.map(({ steam_id }) => steam_id).sort(), + ); + expect(state.fire).toHaveLength(20); + + await Promise.all( + players.flatMap((user) => [ + react(matchId, id, user, "heart"), + react(matchId, id, user, "fire"), + ]), + ); + + expect(await redis.hexists(reactionsKey(matchId), id)).toBe(0); + }); + + it("cancels one player's toggles out in pairs, however they race", async () => { + const user = await author(); + const matchId = randomUUID(); + await seat(matchId, user); + const id = await place(matchId, user); + + await Promise.all( + Array.from({ length: 6 }, () => react(matchId, id, user)), + ); + + expect(await redis.hexists(reactionsKey(matchId), id)).toBe(0); + }); + + // ioredis resends a command whose reply was lost to a reconnect, and a + // toggle that ran twice would undo itself. + it("never runs the same toggle twice", async () => { + const user = await author(); + const matchId = randomUUID(); + const id = await place(matchId, user); + const receipt = `chat_reaction_applied:${randomUUID()}`; + const toggle = () => + redis.eval( + (ChatService as any).TOGGLE_ROOM_REACTION_SCRIPT, + 3, + key(matchId), + reactionsKey(matchId), + receipt, + id, + "heart", + user.steam_id, + "0", + 60_000, + ); + + await expect(toggle()).resolves.toBe( + JSON.stringify({ heart: [user.steam_id] }), + ); + await expect(toggle()).resolves.toBe( + JSON.stringify({ heart: [user.steam_id] }), + ); + + expect(await stored(matchId, id)).toEqual({ heart: [user.steam_id] }); + expect(await redis.pttl(receipt)).toBeGreaterThan(0); + }); + + it("lets a gagged player take a reaction back, but not give one", async () => { + const user = await author(); + const matchId = randomUUID(); + await seat(matchId, user); + const id = await place(matchId, user); + + await react(matchId, id, user, "heart"); + await postgres.query( + `INSERT INTO player_sanctions (player_steam_id, type) + VALUES ($1::bigint, 'gag')`, + [user.steam_id], + ); + + await expect(react(matchId, id, user, "fire")).resolves.toEqual({ + toggled: false, + code: ChatErrorCode.Gagged, + }); + expect(await stored(matchId, id)).toEqual({ heart: [user.steam_id] }); + + await expect(react(matchId, id, user, "heart")).resolves.toEqual({ + toggled: true, + reactions: {}, + }); + expect(await redis.hexists(reactionsKey(matchId), id)).toBe(0); + }); + + it("hands reactions back in the list's order", async () => { + const user = await author(); + const matchId = randomUUID(); + await seat(matchId, user); + const id = await place(matchId, user); + + for (const reaction of ["sad", "fire", "wow", "thumbsup"]) { + await react(matchId, id, user, reaction); + } + + const history = await chat["getMessages"](ChatLobbyType.Match, matchId); + + expect(Object.keys(history[0].reactions)).toEqual([ + "thumbsup", + "fire", + "wow", + "sad", + ]); + }); + + it("stores steam ids as strings, whole", async () => { + const user = await author(); + const matchId = randomUUID(); + await seat(matchId, user); + const id = await place(matchId, user); + + await react(matchId, id, user); + + expect(await redis.hget(reactionsKey(matchId), id)).toBe( + JSON.stringify({ heart: [user.steam_id] }), + ); + }); + + it("gives a message that is not there no reactions", async () => { + const user = await author(); + const matchId = randomUUID(); + await seat(matchId, user); + const expired = await place(matchId, user); + await redis.call("HPEXPIRE", key(matchId), 1, "FIELDS", 1, expired); + await new Promise((resolve) => setTimeout(resolve, 20)); + + for (const id of [randomUUID(), expired]) { + await expect(react(matchId, id, user)).resolves.toEqual({ + toggled: false, + code: ChatErrorCode.NotFound, + }); + } + + expect(await redis.exists(reactionsKey(matchId))).toBe(0); + }); + + it("refuses the ninth toggle in a second, and allows more once it passes", async () => { + const user = await author(); + const matchId = randomUUID(); + await seat(matchId, user); + const id = await place(matchId, user); + + for (let toggle = 0; toggle < ChatService.REACTION_RATE_LIMIT; toggle++) { + await expect(react(matchId, id, user)).resolves.toMatchObject({ + toggled: true, + }); + } + + await expect(react(matchId, id, user)).resolves.toEqual({ + toggled: false, + code: ChatErrorCode.RateLimited, + }); + + const ttl = await redis.pttl(`chat:reaction-rate:${user.steam_id}`); + expect(ttl).toBeGreaterThan(0); + expect(ttl).toBeLessThanOrEqual(1_000); + + await new Promise((resolve) => setTimeout(resolve, ttl + 50)); + + await expect(react(matchId, id, user)).resolves.toMatchObject({ + toggled: true, + }); + }); + + it("takes a message's reactions with it when it is deleted", async () => { + const user = await author(); + const matchId = randomUUID(); + await seat(matchId, user); + const id = await place(matchId, user); + const kept = await place(matchId, user); + + await react(matchId, id, user); + await react(matchId, kept, user); + + await expect( + chat.deleteMessage(ChatLobbyType.Match, matchId, id, user), + ).resolves.toEqual({ deleted: true }); + + expect(await redis.hexists(reactionsKey(matchId), id)).toBe(0); + expect(await stored(matchId, kept)).toEqual({ heart: [user.steam_id] }); + }); + + it("puts each message's reactions in the room's history", async () => { + const user = await author(); + const matchId = randomUUID(); + await seat(matchId, user); + const reacted = await place(matchId, user); + const quiet = await place(matchId, user); + + await react(matchId, reacted, user, "wow"); + + const history = await chat["getMessages"](ChatLobbyType.Match, matchId); + + expect( + Object.fromEntries(history.map(({ id, reactions }) => [id, reactions])), + ).toEqual({ + [reacted]: { wow: [user.steam_id] }, + [quiet]: {}, + }); + }); + + it("carries reactions from a draft into its match in the same step", async () => { + const organizer = { ...(await author()), role: "match_organizer" }; + const draftId = randomUUID(); + const matchId = randomUUID(); + const id = randomUUID(); + await seat(draftId, organizer, "draft"); + await redis.hset(`chat_draft_${draftId}`, id, written(id, organizer)); + + await react(draftId, id, organizer, "sad", ChatLobbyType.Draft); + + await chat.migrateLobbyMessages( + ChatLobbyType.Draft, + draftId, + ChatLobbyType.Match, + matchId, + ); + + expect(await stored(matchId, id)).toEqual({ + sad: [organizer.steam_id], + }); + expect(await expiresAt(matchId, id)).toBeGreaterThan(Date.now()); + expect(await reactionsExpireAt(matchId, id)).toBe( + await expiresAt(matchId, id), + ); + expect(await redis.exists(reactionsKey(draftId, "draft"))).toBe(0); + expect(await chat["getMessages"](ChatLobbyType.Match, matchId)).toEqual([ + expect.objectContaining({ + id, + reactions: { sad: [organizer.steam_id] }, + }), + ]); + + await seat(draftId, organizer, "draft"); + + await expect( + react(draftId, id, organizer, "sad", ChatLobbyType.Draft), + ).resolves.toEqual({ toggled: false, code: ChatErrorCode.NotFound }); + expect(await redis.exists(reactionsKey(draftId, "draft"))).toBe(0); + }); + + it("leaves no reactions behind when the match room keeps nothing", async () => { + const organizer = { ...(await author()), role: "match_organizer" }; + const draftId = randomUUID(); + const matchId = randomUUID(); + const id = randomUUID(); + await seat(draftId, organizer, "draft"); + await redis.hset(`chat_draft_${draftId}`, id, written(id, organizer)); + await react(draftId, id, organizer, "sad", ChatLobbyType.Draft); + + chat.updateChatMessageTTL(ChatLobbyType.Match, 0); + + try { + await chat.migrateLobbyMessages( + ChatLobbyType.Draft, + draftId, + ChatLobbyType.Match, + matchId, + ); + } finally { + chat.updateChatMessageTTL(ChatLobbyType.Match, 60 * 60); + } + + expect(await redis.exists(key(matchId))).toBe(0); + expect(await redis.exists(reactionsKey(matchId))).toBe(0); + expect(await redis.exists(reactionsKey(draftId, "draft"))).toBe(0); + }); + }); + describe("the bell's preview", () => { const row = async ( steamId: string,