From b153d7fd3df976f1a489bf442b4bd7c2d95ec1cb Mon Sep 17 00:00:00 2001 From: Luke Policinski Date: Mon, 28 Sep 2026 17:05:42 -0400 Subject: [PATCH 1/3] feature: chat message reactions Players can react to a chat message with one of six reactions (thumbsup, heart, laugh, fire, wow, sad) and take it back by sending it again. - lobby:react { type, id, messageId, reaction, requestId? } toggles the sender's reaction. The room gets lobby:::reaction { id, reactions } with the message's whole reaction state, never a delta. chat:ack / chat:error carry action "react". - Held to the same rules as sending: room access (canPostIn), and a gag blocks reacting in group rooms. Never a notification, never relayed to the game. At most 8 toggles a second per player (rate_limited), since every toggle is a broadcast to the room. - Rooms: a separate chat_reactions__ hash, field per message, so reactions never contend with an edit's compare-and-set. One Lua script toggles, refuses a message that is gone, and gives the reactions the message's own expiry. Deleting a message clears its reactions, and the draft-to-match move carries them in the same step. - DMs: direct_message_reactions, toggled in one statement scoped to the conversation; cascades with the message, the player and retention. - History (room and DM) and live messages carry reactions per message. --- .../public_direct_message_reactions.yaml | 3 + .../databases/default/tables/tables.yaml | 1 + .../down.sql | 1 + .../up.sql | 10 + src/chat/chat.gateway.spec.ts | 147 +++++ src/chat/chat.gateway.ts | 50 ++ src/chat/chat.service.spec.ts | 563 +++++++++++++++++- src/chat/chat.service.ts | 312 +++++++++- src/chat/enums/ChatErrorCode.ts | 1 + src/chat/types/ChatAction.ts | 2 +- src/chat/types/ChatMessage.ts | 4 + src/chat/types/ChatReactResult.ts | 6 + src/chat/types/ChatReactions.ts | 3 + test/chat-direct-messages.spec.ts | 176 ++++++ test/chat-redis-actions.spec.ts | 313 ++++++++++ 15 files changed, 1570 insertions(+), 22 deletions(-) create mode 100644 hasura/metadata/databases/default/tables/public_direct_message_reactions.yaml create mode 100644 hasura/migrations/default/1888000000300_direct_message_reactions/down.sql create mode 100644 hasura/migrations/default/1888000000300_direct_message_reactions/up.sql create mode 100644 src/chat/types/ChatReactResult.ts create mode 100644 src/chat/types/ChatReactions.ts 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..fd8678e9 --- /dev/null +++ b/hasura/migrations/default/1888000000300_direct_message_reactions/up.sql @@ -0,0 +1,10 @@ +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) +); 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..5ffb7a12 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 & { code?: string }) | undefined; // The one direct message the fake database holds, if a test put one there. let directMessage: | { @@ -70,6 +73,20 @@ describe("ChatService direct messages", () => { return [{ deleted: audited }]; } + 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] && + directMessage.roomId === bindings[1] + ? [{ reactions: directReactions }] + : []; + } + if (sql.includes("INSERT INTO public.chat_message_edits")) { const id = `edit-audit-${editAuditIds.length + 1}`; editAuditIds.push(id); @@ -328,6 +345,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 +366,8 @@ describe("ChatService direct messages", () => { gagged = false; audited = false; editAuditIds = []; + directReactions = null; + directReactionFailure = undefined; rcon.send.mockResolvedValue(undefined); rcon.connect.mockResolvedValue(rcon); @@ -2195,6 +2215,547 @@ 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; + + 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 = {}; + 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, messageId, reaction, steamId] = args; + + if (!hashes[roomKey]?.[messageId]) { + return null; + } + + const state = JSON.parse(hashes[reactionsKey]?.[messageId] ?? "{}"); + const holders: string[] = state[reaction] ?? []; + 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 and who reacted", async () => { + await react("laugh"); + + expect(toggles()[0].slice(1)).toEqual([ + 2, + ROOM, + REACTIONS, + MESSAGE_ID, + "laugh", + ME, + ]); + }); + + 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 reacting in a group room", async () => { + gagged = true; + + await expect(react()).resolves.toEqual({ + toggled: false, + code: ChatErrorCode.Gagged, + }); + expect(toggles()).toHaveLength(0); + }); + + 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"); + }); + + it("moves reactions with a draft's messages into its match", async () => { + await react("fire"); + 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("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 toggle = queries.find(({ sql }) => + sql.includes("INSERT INTO public.direct_message_reactions"), + ); + + 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 = undefined; + + await expect(reactDirect()).resolves.toEqual({ + toggled: false, + code: ChatErrorCode.NotFound, + }); + await flush(); + + expect(broadcasts(`lobby:direct:${room}:reaction`)).toEqual([]); + }); + + it("answers not_found when the message is deleted under the toggle", async () => { + directReactionFailure = Object.assign( + new Error("violates foreign key constraint"), + { code: "23503" }, + ); + + await expect(reactDirect()).resolves.toEqual({ + toggled: false, + code: ChatErrorCode.NotFound, + }); + }); + + 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..fee6b68d 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,113 @@ 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. + private static readonly TOGGLE_ROOM_REACTION_SCRIPT = ` + 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 + table.insert(reactors, ARGV[3]) + end + if #reactors == 0 then + reactions[ARGV[2]] = nil + else + reactions[ARGV[2]] = reactors + end + 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 + `; + + // 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 +579,40 @@ 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: reactions[messageId] + ? JSON.parse(reactions[messageId]) + : {}, + }), + ) .sort( (a, b) => new Date(a.timestamp).getTime() - new Date(b.timestamp).getTime(), ); } + 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 +810,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 +912,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 +1016,145 @@ 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 }; + } + + if (type !== ChatLobbyType.Direct && (await this.isGagged(user.steam_id))) { + return { toggled: false, code: ChatErrorCode.Gagged }; + } + + const reactions = + type === ChatLobbyType.Direct + ? await this.toggleDirectReaction(id, messageId, reaction, user) + : await this.toggleRoomReaction(type, id, messageId, reaction, user); + + 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, + ): Promise { + const state = await this.redis.eval( + ChatService.TOGGLE_ROOM_REACTION_SCRIPT, + 2, + `chat_${type}_${id}`, + ChatService.reactionsKey(type, id), + messageId, + reaction, + String(user.steam_id), + ); + + if (typeof state !== "string") { + return null; + } + + return JSON.parse(state) as ChatReactions; + } + + private async toggleDirectReaction( + roomId: string, + messageId: string, + reaction: string, + user: User, + ): Promise { + try { + await this.postgres.query( + `WITH message AS ( + SELECT id FROM public.direct_messages + WHERE id = $1::uuid AND room_id = $4 + ), removed AS ( + DELETE FROM public.direct_message_reactions r + USING message + WHERE r.message_id = message.id + AND r.steam_id = $2::bigint + AND r.reaction = $3 + RETURNING 1 + ) + INSERT INTO public.direct_message_reactions + (message_id, steam_id, reaction) + SELECT message.id, $2::bigint, $3 + FROM message + WHERE NOT EXISTS (SELECT 1 FROM removed) + ON CONFLICT DO NOTHING`, + [messageId, String(user.steam_id), reaction, roomId], + ); + } catch (error) { + // A message deleted after the statement found it fails the foreign + // key, rather than matching nothing. + if (error?.code === "23503") { + return null; + } + + throw error; + } + + const [row] = await this.postgres.query< + Array<{ reactions: ChatReactions | null }> + >( + `SELECT reactions.reactions + FROM public.direct_messages dm + ${ChatService.DIRECT_MESSAGE_REACTIONS} + WHERE dm.id = $1::uuid AND dm.room_id = $2`, + [messageId, roomId], + ); + + if (!row) { + return null; + } + + return row.reactions ?? {}; + } + private async editRoomMessage( type: ChatLobbyType, id: string, @@ -1887,6 +2160,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 +2169,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 +2191,7 @@ export class ChatService { ...(row.edited_at ? { edited_at: new Date(row.edited_at).toISOString() } : {}), + reactions: row.reactions ?? {}, from: { role: row.role, name: row.name, @@ -2079,6 +2356,7 @@ export class ChatService { | "chat" | "edited" | "deleted" + | "reaction" | "list" | "messages" | "joined" @@ -2407,13 +2685,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 +2701,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..105ad7e7 100644 --- a/test/chat-direct-messages.spec.ts +++ b/test/chat-direct-messages.spec.ts @@ -653,6 +653,182 @@ 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("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"); + + 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..7c82667e 100644 --- a/test/chat-redis-actions.spec.ts +++ b/test/chat-redis-actions.spec.ts @@ -687,6 +687,319 @@ 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("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); + }); + + 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, From d09cdb8e081e47626f8354e36d59bbe847b756cc Mon Sep 17 00:00:00 2001 From: Luke Policinski Date: Mon, 28 Sep 2026 17:22:32 -0400 Subject: [PATCH 2/3] =?UTF-8?q?bug:=20chat=20reactions=20review=20fixes=20?= =?UTF-8?q?=E2=80=94=20resend=20receipt,=20gagged=20un-react,=20serialized?= =?UTF-8?q?=20DM=20toggles?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - The room toggle script takes a receipt key per toggle: ioredis resends a command whose reply was lost to a reconnect, and a toggle run twice undid itself. A resend now answers with the current state instead. - A gagged player can still take back a reaction they gave, the way #442 lets them delete their own message; only adding one is refused. - DM toggles lock the message row first (FOR NO KEY UPDATE) and toggle in a fresh statement, so racing toggles on an absent row cancel out as they do in rooms instead of both inserting, and a delete can no longer land between finding the message and writing the reaction. - Reactions are always handed out in ChatService.REACTIONS order; the Lua table and jsonb each had their own. - Index direct_message_reactions (steam_id) for the players cascade. --- .../up.sql | 3 + src/chat/chat.service.spec.ts | 134 ++++++++++++++--- src/chat/chat.service.ts | 142 ++++++++++++------ test/chat-direct-messages.spec.ts | 54 +++++++ test/chat-redis-actions.spec.ts | 100 ++++++++++++ 5 files changed, 362 insertions(+), 71 deletions(-) diff --git a/hasura/migrations/default/1888000000300_direct_message_reactions/up.sql b/hasura/migrations/default/1888000000300_direct_message_reactions/up.sql index fd8678e9..8cc66450 100644 --- a/hasura/migrations/default/1888000000300_direct_message_reactions/up.sql +++ b/hasura/migrations/default/1888000000300_direct_message_reactions/up.sql @@ -8,3 +8,6 @@ CREATE TABLE IF NOT EXISTS public.direct_message_reactions ( 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.service.spec.ts b/src/chat/chat.service.spec.ts index 5ffb7a12..ab2e8f60 100644 --- a/src/chat/chat.service.spec.ts +++ b/src/chat/chat.service.spec.ts @@ -41,7 +41,7 @@ describe("ChatService direct messages", () => { let audited: boolean; let editAuditIds: string[]; let directReactions: Record | null; - let directReactionFailure: (Error & { code?: string }) | undefined; + let directReactionFailure: Error | undefined; // The one direct message the fake database holds, if a test put one there. let directMessage: | { @@ -73,6 +73,13 @@ 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; @@ -81,8 +88,7 @@ describe("ChatService direct messages", () => { } if (sql.includes("SELECT reactions.reactions")) { - return directMessage?.id === bindings[0] && - directMessage.roomId === bindings[1] + return directMessage?.id === bindings[0] ? [{ reactions: directReactions }] : []; } @@ -144,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(), @@ -2222,6 +2239,7 @@ describe("ChatService direct messages", () => { let hashes: Record>; let rates: Record; + let receipts: Set; const as = (steamId: string, overrides: Record = {}) => ({ @@ -2251,6 +2269,7 @@ describe("ChatService direct messages", () => { beforeEach(() => { rates = {}; + receipts = new Set(); hashes = { [ROOM]: { [MESSAGE_ID]: JSON.stringify({ @@ -2290,7 +2309,19 @@ describe("ChatService direct messages", () => { } if (script.includes("cjson")) { - const [roomKey, reactionsKey, messageId, reaction, steamId] = args; + const [ + roomKey, + reactionsKey, + receipt, + messageId, + reaction, + steamId, + removeOnly, + ] = args; + + if (receipts.has(receipt)) { + return hashes[reactionsKey]?.[messageId] ?? "{}"; + } if (!hashes[roomKey]?.[messageId]) { return null; @@ -2298,6 +2329,12 @@ describe("ChatService direct messages", () => { 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]; @@ -2344,16 +2381,36 @@ describe("ChatService direct messages", () => { }); }); - it("hands the toggle both hashes and who reacted", async () => { + it("hands the toggle both hashes, a receipt of its own and who reacted", async () => { + await react("laugh"); await react("laugh"); - expect(toggles()[0].slice(1)).toEqual([ - 2, + 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", ]); }); @@ -2432,14 +2489,28 @@ describe("ChatService direct messages", () => { expect(toggles()).toHaveLength(0); }); - it("keeps a gagged player from reacting in a group room", async () => { + 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, }); - expect(toggles()).toHaveLength(0); + 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 () => { @@ -2572,8 +2643,10 @@ describe("ChatService direct messages", () => { expect(stored).not.toHaveProperty("reactions"); }); - it("moves reactions with a draft's messages into its match", async () => { - await react("fire"); + // 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( @@ -2634,17 +2707,24 @@ describe("ChatService direct messages", () => { ); }); - it("toggles in one statement scoped to the conversation, then reads the state", async () => { + 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 toggle = queries.find(({ sql }) => + 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", @@ -2689,7 +2769,10 @@ describe("ChatService direct messages", () => { }); it("answers not_found for a message that is not in the conversation", async () => { - directMessage = undefined; + directMessage = { + ...directMessage, + roomId: directRoomId(FRIEND, STRANGER), + }; await expect(reactDirect()).resolves.toEqual({ toggled: false, @@ -2697,19 +2780,24 @@ describe("ChatService direct messages", () => { }); await flush(); + expect( + queries.some(({ sql }) => + sql.includes("INSERT INTO public.direct_message_reactions"), + ), + ).toBe(false); expect(broadcasts(`lobby:direct:${room}:reaction`)).toEqual([]); }); - it("answers not_found when the message is deleted under the toggle", async () => { - directReactionFailure = Object.assign( - new Error("violates foreign key constraint"), - { code: "23503" }, - ); + it("orders reactions as the list does", async () => { + directReactions = { sad: [FRIEND], heart: [ME], thumbsup: [FRIEND] }; - await expect(reactDirect()).resolves.toEqual({ - toggled: false, - code: ChatErrorCode.NotFound, - }); + const result = await reactDirect(); + + expect(Object.keys(result.toggled ? result.reactions : {})).toEqual([ + "thumbsup", + "heart", + "sad", + ]); }); it("lets any other failure through", async () => { diff --git a/src/chat/chat.service.ts b/src/chat/chat.service.ts index fee6b68d..ca22231f 100644 --- a/src/chat/chat.service.ts +++ b/src/chat/chat.service.ts @@ -144,7 +144,15 @@ export class ChatService { // 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 @@ -163,6 +171,9 @@ export class ChatService { end end if not removed then + if ARGV[4] == '1' then + return 0 + end table.insert(reactors, ARGV[3]) end if #reactors == 0 then @@ -170,6 +181,7 @@ export class ChatService { 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 '{}' @@ -183,6 +195,8 @@ export class ChatService { 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 ( @@ -598,9 +612,9 @@ export class ChatService { .map( ([messageId, value]): ChatMessage => ({ ...(JSON.parse(value) as ChatMessage), - reactions: reactions[messageId] - ? JSON.parse(reactions[messageId]) - : {}, + reactions: ChatService.orderedReactions( + reactions[messageId] ? JSON.parse(reactions[messageId]) : null, + ), }), ) .sort( @@ -609,6 +623,18 @@ export class ChatService { ); } + // 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}`; } @@ -1048,14 +1074,23 @@ export class ChatService { return { toggled: false, code: ChatErrorCode.NotAllowed }; } - if (type !== ChatLobbyType.Direct && (await this.isGagged(user.steam_id))) { - return { toggled: false, code: ChatErrorCode.Gagged }; - } - const reactions = type === ChatLobbyType.Direct ? await this.toggleDirectReaction(id, messageId, reaction, user) - : await this.toggleRoomReaction(type, 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 }; @@ -1083,22 +1118,30 @@ export class ChatService { messageId: string, reaction: string, user: User, - ): Promise { + removeOnly: boolean, + ): Promise { const state = await this.redis.eval( ChatService.TOGGLE_ROOM_REACTION_SCRIPT, - 2, + 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 JSON.parse(state) as ChatReactions; + return ChatService.orderedReactions(JSON.parse(state)); } private async toggleDirectReaction( @@ -1107,52 +1150,55 @@ export class ChatService { reaction: string, user: User, ): Promise { - try { - await this.postgres.query( - `WITH message AS ( - SELECT id FROM public.direct_messages - WHERE id = $1::uuid AND room_id = $4 - ), removed AS ( - DELETE FROM public.direct_message_reactions r - USING message - WHERE r.message_id = message.id - AND r.steam_id = $2::bigint - AND r.reaction = $3 + 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 message.id, $2::bigint, $3 - FROM message + 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], ); - } catch (error) { - // A message deleted after the statement found it fails the foreign - // key, rather than matching nothing. - if (error?.code === "23503") { - return null; - } - throw error; - } - - const [row] = await this.postgres.query< - Array<{ reactions: ChatReactions | null }> - >( - `SELECT reactions.reactions - FROM public.direct_messages dm - ${ChatService.DIRECT_MESSAGE_REACTIONS} - WHERE dm.id = $1::uuid AND dm.room_id = $2`, - [messageId, roomId], - ); - - if (!row) { - return null; - } + 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 row.reactions ?? {}; + return ChatService.orderedReactions(row?.reactions); + }); } private async editRoomMessage( @@ -2191,7 +2237,7 @@ export class ChatService { ...(row.edited_at ? { edited_at: new Date(row.edited_at).toISOString() } : {}), - reactions: row.reactions ?? {}, + reactions: ChatService.orderedReactions(row.reactions), from: { role: row.role, name: row.name, diff --git a/test/chat-direct-messages.spec.ts b/test/chat-direct-messages.spec.ts index 105ad7e7..5c68ee1f 100644 --- a/test/chat-direct-messages.spec.ts +++ b/test/chat-direct-messages.spec.ts @@ -722,6 +722,49 @@ describe("direct messages (SQL-driven)", () => { 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); + 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(); @@ -820,6 +863,17 @@ describe("direct messages (SQL-driven)", () => { 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(); diff --git a/test/chat-redis-actions.spec.ts b/test/chat-redis-actions.spec.ts index 7c82667e..ad0a5e02 100644 --- a/test/chat-redis-actions.spec.ts +++ b/test/chat-redis-actions.spec.ts @@ -779,6 +779,28 @@ describe("chat edits and self deletes (SQL-driven)", () => { 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(); @@ -836,6 +858,84 @@ describe("chat edits and self deletes (SQL-driven)", () => { 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(); From 8754dee0c6b33c3bc225f08feae7c4474b6e05f1 Mon Sep 17 00:00:00 2001 From: Luke Policinski Date: Mon, 28 Sep 2026 20:37:33 -0400 Subject: [PATCH 3/3] test: prove racing DM reaction toggles are serialized One round of the race lost it without the row lock only about three runs in four, so the regression could pass. Eight rounds fail every time without the lock. --- test/chat-direct-messages.spec.ts | 19 ++++++++++++------- 1 file changed, 12 insertions(+), 7 deletions(-) diff --git a/test/chat-direct-messages.spec.ts b/test/chat-direct-messages.spec.ts index 5c68ee1f..f179ba43 100644 --- a/test/chat-direct-messages.spec.ts +++ b/test/chat-direct-messages.spec.ts @@ -726,15 +726,20 @@ describe("direct messages (SQL-driven)", () => { const me = await fx.player(); const friend = await fx.player(); const room = directRoomId(me, friend); - 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)), - ]); + // 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 }]); + expect(await rows()).toEqual([{ message_id: odd }]); + } }); it("counts both parties toggling every reaction at once exactly", async () => {