From 912c9d9ea2083cc73baf0144646bf82f35c713bb Mon Sep 17 00:00:00 2001 From: Luke Policinski Date: Mon, 28 Sep 2026 16:22:35 -0400 Subject: [PATCH 1/5] feature: edit or delete your own chat messages for 10 minutes The author of a website message (source "web") can edit or delete it for 10 minutes (ChatService.SELF_SERVICE_WINDOW_MS, 5s of pod clock skew allowed). Nobody edits another player's message, administrators included; moderator deletes at any age stay as they were. - lobby:edit { type, id, messageId, message, requestId? } broadcasts lobby:::edited { id, message, edited_at }, acks/errors with action "edit", and is never relayed to the game server - room edits are one compare-and-set Lua script against the value that was read and checked: HSET drops a field's TTL, so the exact HPEXPIRETIME is put back with HPEXPIREAT, and an expired or deleted field is never written back - a gag blocks editing in group rooms but not deleting your own message; an author's own delete in a group room is audited like a moderator's - DMs: direct_messages.edited_at (migration 1888000000200); edit and delete are single statements guarded on author, room and created_at on the database clock; not audited; history returns edited_at - unread, live bell rows get the edited preview; rows a delete already retracted stay blank, and an edit or delete that lands while a message's rows are still being written is applied once they are - new chat:error code window_closed --- .../down.sql | 2 + .../up.sql | 2 + src/chat/chat.gateway.spec.ts | 147 ++++ src/chat/chat.gateway.ts | 54 ++ src/chat/chat.service.spec.ts | 706 +++++++++++++++++- src/chat/chat.service.ts | 447 ++++++++++- src/chat/enums/ChatErrorCode.ts | 1 + src/chat/types/ChatAction.ts | 2 +- src/chat/types/ChatEditResult.ts | 5 + src/chat/types/ChatMessage.ts | 2 + src/notifications/notifications.service.ts | 14 + test/chat-direct-messages.spec.ts | 264 ++++++- test/chat-redis-actions.spec.ts | 469 ++++++++++++ 13 files changed, 2085 insertions(+), 30 deletions(-) create mode 100644 hasura/migrations/default/1888000000200_direct_messages_edited_at/down.sql create mode 100644 hasura/migrations/default/1888000000200_direct_messages_edited_at/up.sql create mode 100644 src/chat/types/ChatEditResult.ts create mode 100644 test/chat-redis-actions.spec.ts diff --git a/hasura/migrations/default/1888000000200_direct_messages_edited_at/down.sql b/hasura/migrations/default/1888000000200_direct_messages_edited_at/down.sql new file mode 100644 index 00000000..3a460907 --- /dev/null +++ b/hasura/migrations/default/1888000000200_direct_messages_edited_at/down.sql @@ -0,0 +1,2 @@ +ALTER TABLE public.direct_messages + DROP COLUMN IF EXISTS edited_at; diff --git a/hasura/migrations/default/1888000000200_direct_messages_edited_at/up.sql b/hasura/migrations/default/1888000000200_direct_messages_edited_at/up.sql new file mode 100644 index 00000000..d342e0fc --- /dev/null +++ b/hasura/migrations/default/1888000000200_direct_messages_edited_at/up.sql @@ -0,0 +1,2 @@ +ALTER TABLE public.direct_messages + ADD COLUMN IF NOT EXISTS edited_at timestamptz; diff --git a/src/chat/chat.gateway.spec.ts b/src/chat/chat.gateway.spec.ts index 730692d8..bb5c3a87 100644 --- a/src/chat/chat.gateway.spec.ts +++ b/src/chat/chat.gateway.spec.ts @@ -461,3 +461,150 @@ describe("ChatGateway lobby:delete", () => { ]); }); }); + +describe("ChatGateway lobby:edit", () => { + const MESSAGE_ID = "3f0c1d2e-4b5a-4c6d-8e7f-9a0b1c2d3e4f"; + + let chat: { editMessage: 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 edit = (overrides: Record = {}) => ({ + id: "m-1", + type: ChatLobbyType.Match, + messageId: MESSAGE_ID, + message: "fixed", + ...overrides, + }); + + beforeEach(() => { + chat = { + editMessage: jest.fn().mockResolvedValue({ + edited: true, + message: "fixed", + edited_at: "2026-01-01T00:00:00.000Z", + }), + 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.editMessage(edit({ requestId: "r-1" }) as any, socket); + + expect(chat.editMessage).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 }], + ["a message that is not a string", { message: 5 }], + ["a message that is only whitespace", { message: " \n " }], + ])("ignores %s", async (_, overrides) => { + const socket = client(); + + await gateway.editMessage(edit(overrides) as any, socket); + + expect(chat.editMessage).not.toHaveBeenCalled(); + expect(socket.send).not.toHaveBeenCalled(); + }); + + it("ignores a missing payload", async () => { + await gateway.editMessage(undefined as any, client()); + + expect(chat.editMessage).not.toHaveBeenCalled(); + }); + + it("asks the service to edit as the signed in player, with the trimmed text", async () => { + await gateway.editMessage(edit({ message: " fixed " }) as any, client()); + + expect(chat.editMessage).toHaveBeenCalledWith( + ChatLobbyType.Match, + "m-1", + MESSAGE_ID, + expect.objectContaining({ steam_id: "1" }), + "fixed", + ); + }); + + it("refuses an edit over the limit without asking the service", async () => { + const socket = client(); + + await gateway.editMessage( + edit({ + message: "a".repeat(ChatService.MAX_MESSAGE_LENGTH + 1), + requestId: "r-1", + }) as any, + socket, + ); + + expect(chat.editMessage).not.toHaveBeenCalled(); + expect(sent(socket)).toEqual([ + { + event: "chat:error", + data: { + code: ChatErrorCode.TooLong, + action: "edit", + max: 2000, + requestId: "r-1", + }, + }, + ]); + }); + + it("acks an edit under the requestId it came with", async () => { + const socket = client(); + + await gateway.editMessage(edit({ requestId: "r-2" }) as any, socket); + + expect(sent(socket)).toEqual([ + { + event: "chat:ack", + data: { requestId: "r-2", messageId: MESSAGE_ID, action: "edit" }, + }, + ]); + }); + + it("stays quiet for an edit without a requestId", async () => { + const socket = client(); + + await gateway.editMessage(edit() as any, socket); + + expect(socket.send).not.toHaveBeenCalled(); + }); + + it.each([ + ChatErrorCode.NotAllowed, + ChatErrorCode.NotFound, + ChatErrorCode.WindowClosed, + ChatErrorCode.Gagged, + ])("reports %s under the requestId it came with", async (code) => { + chat.editMessage.mockResolvedValue({ edited: false, code }); + const socket = client(); + + await gateway.editMessage(edit({ requestId: "r-3" }) as any, socket); + + expect(sent(socket)).toEqual([ + { + event: "chat:error", + data: { code, action: "edit", requestId: "r-3" }, + }, + ]); + }); + + it("never relays an edit to the game server", async () => { + await gateway.editMessage(edit() as any, client()); + + expect(chat.editMessage).toHaveBeenCalled(); + expect(chat.sendChatToServer).not.toHaveBeenCalled(); + }); +}); diff --git a/src/chat/chat.gateway.ts b/src/chat/chat.gateway.ts index 6a9e58c3..59df2690 100644 --- a/src/chat/chat.gateway.ts +++ b/src/chat/chat.gateway.ts @@ -195,6 +195,60 @@ export class ChatGateway { } } + @SubscribeMessage("lobby:edit") + async editMessage( + @MessageBody() + data: { + id: string; + type: ChatLobbyType; + messageId: string; + message: 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; + + const parsed = ChatService.messageText(data.message); + + if ("error" in parsed) { + if (parsed.error === ChatErrorCode.TooLong) { + this.sendError(client, "edit", parsed.error, requestId); + } + return; + } + + const result = await this.chat.editMessage( + data.type, + data.id, + data.messageId, + client.user, + parsed.text, + ); + + if (result.edited === false) { + this.sendError(client, "edit", result.code, requestId); + return; + } + + if (requestId) { + this.sendAck(client, "edit", 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 196753cd..91339fae 100644 --- a/src/chat/chat.service.spec.ts +++ b/src/chat/chat.service.spec.ts @@ -38,6 +38,25 @@ describe("ChatService direct messages", () => { let queries: Array<{ sql: string; bindings: any[] }>; let gagged: boolean; let audited: boolean; + // The one direct message the fake database holds, if a test put one there. + let directMessage: + | { + id: string; + roomId: string; + author: string; + message: string; + open: boolean; + editedAt: Date | null; + } + | undefined; + + const ownsDirectMessage = (bindings: any[]) => + directMessage !== undefined && + directMessage.id === bindings[0] && + directMessage.roomId === bindings[1] && + directMessage.author === bindings[2] && + directMessage.open; + const postgres = { query: jest.fn(async (sql: string, bindings: any[]): Promise => { queries.push({ sql, bindings }); @@ -50,6 +69,53 @@ describe("ChatService direct messages", () => { return [{ deleted: audited }]; } + if (sql.includes("AS open")) { + return directMessage?.id === bindings[0] && + directMessage.roomId === bindings[1] + ? [{ author: directMessage.author, open: directMessage.open }] + : []; + } + + if (sql.includes("UPDATE public.direct_messages")) { + if (!ownsDirectMessage(bindings)) { + return []; + } + + directMessage.message = bindings[4]; + directMessage.editedAt = new Date("2026-01-01T00:00:00.000Z"); + + return [ + { + message: directMessage.message, + edited_at: directMessage.editedAt, + }, + ]; + } + + if (sql.includes("DELETE FROM public.direct_messages")) { + if (!ownsDirectMessage(bindings)) { + return []; + } + + const { id } = directMessage; + directMessage = undefined; + + return [{ id }]; + } + + if ( + sql.includes("SELECT message, edited_at FROM public.direct_messages") + ) { + return directMessage?.id === bindings[0] + ? [ + { + message: directMessage.message, + edited_at: directMessage.editedAt, + }, + ] + : []; + } + return []; }), }; @@ -59,6 +125,7 @@ describe("ChatService direct messages", () => { collapseOlderUnread: jest.fn(), markConversationRead: jest.fn(), retractChatMessage: jest.fn().mockResolvedValue(undefined), + updateChatMessagePreview: jest.fn().mockResolvedValue(undefined), }; const client = (steamId: string) => @@ -255,7 +322,10 @@ describe("ChatService direct messages", () => { redis.hget.mockResolvedValue(null); redis.hgetall.mockResolvedValue({}); redis.get.mockResolvedValue(null); + redis.eval.mockResolvedValue([1, 1]); notifications.retractChatMessage.mockResolvedValue(undefined); + notifications.updateChatMessagePreview.mockResolvedValue(undefined); + directMessage = undefined; acceptedFriendships = [[ME, FRIEND]]; myMatches = ["m-1"]; otherMatches = ["mm-1"]; @@ -1178,8 +1248,16 @@ describe("ChatService direct messages", () => { expect(redis.hdel).not.toHaveBeenCalled(); }); - it("refuses anyone in a direct conversation", async () => { + it("does not let an administrator moderate a direct conversation", async () => { role = "administrator"; + directMessage = { + id: MESSAGE_ID, + roomId: directRoomId(ME, FRIEND), + author: FRIEND, + message: "something awful", + open: true, + editedAt: null, + }; await expect( service.deleteMessage( @@ -1190,8 +1268,8 @@ describe("ChatService direct messages", () => { ), ).resolves.toEqual({ deleted: false, code: ChatErrorCode.NotAllowed }); - expect(redis.hget).not.toHaveBeenCalled(); - expect(queries).toHaveLength(0); + expect(directMessage).toBeDefined(); + expect(audits()).toHaveLength(0); }); it("answers not_found for a message that is not there", async () => { @@ -1291,6 +1369,628 @@ describe("ChatService direct messages", () => { }); }); + describe("editing and deleting your own messages", () => { + const MESSAGE_ID = "5a1b2c3d-4e5f-4a6b-8c7d-9e0f1a2b3c4d"; + const WINDOW = ChatService.SELF_SERVICE_WINDOW_MS; + + let stored: Record>; + + const ago = (ms: number) => new Date(Date.now() - ms).toISOString(); + + const store = (message: Record = {}, id = "m-1") => { + stored[`chat_match_${id}`] = { + [MESSAGE_ID]: JSON.stringify({ + id: MESSAGE_ID, + message: "typo", + timestamp: ago(60_000), + source: "web", + from: { + role: "user", + name: "Me", + steam_id: ME, + avatar_url: null, + profile_url: "https://steamcommunity.com/id/me/", + }, + ...message, + }), + }; + }; + + const current = (id = "m-1") => + JSON.parse(stored[`chat_match_${id}`]?.[MESSAGE_ID] ?? "null"); + + const me = (overrides: Record = {}) => + ({ steam_id: ME, name: "Me", role: "user", ...overrides }) as any; + + const edit = (text = "fixed", user = me(), id = "m-1") => + service.editMessage(ChatLobbyType.Match, id, MESSAGE_ID, user, text); + + const selfDelete = (user = me(), id = "m-1") => + service.deleteMessage(ChatLobbyType.Match, id, MESSAGE_ID, user); + + const broadcasts = (event: string) => + redis.publish.mock.calls + .map(([, payload]) => JSON.parse(payload)) + .filter((published) => published.event === event); + + const audits = () => + queries.filter(({ sql }) => + sql.includes("INSERT INTO public.chat_message_deletions"), + ); + + const flush = () => new Promise((resolve) => setImmediate(resolve)); + + beforeEach(() => { + stored = {}; + redis.hget.mockImplementation( + async (key: string, field: string) => stored[key]?.[field] ?? null, + ); + redis.hgetall.mockImplementation(async (key: string) => + key === "chat:match:m-1" + ? { [FRIEND]: JSON.stringify({ user: { steam_id: FRIEND } }) } + : {}, + ); + redis.hdel.mockImplementation(async (key: string, field: string) => { + delete stored[key]?.[field]; + return 1; + }); + redis.eval.mockImplementation( + async ( + script: string, + _keys: number, + key: string, + field: string, + expected: string, + next: string, + ) => { + if (!script.includes("HPEXPIRETIME")) { + return [1, 1]; + } + + if (stored[key]?.[field] !== expected) { + return 0; + } + + stored[key][field] = next; + return 1; + }, + ); + }); + + describe("in a room", () => { + it("lets the author change what they wrote within the window", async () => { + store(); + const before = current(); + + const result = await edit(" fixed "); + + expect(result).toEqual({ + edited: true, + message: "fixed", + edited_at: expect.any(String), + }); + expect(current()).toEqual({ + ...before, + message: "fixed", + edited_at: result.edited ? result.edited_at : undefined, + }); + }); + + it("tells the room what the message says now", async () => { + store(); + + const result = await edit(); + await flush(); + + expect(broadcasts("lobby:match:m-1:edited")).toEqual([ + { + steamId: FRIEND, + event: "lobby:match:m-1:edited", + data: { + id: MESSAGE_ID, + message: "fixed", + edited_at: result.edited ? result.edited_at : undefined, + }, + }, + ]); + expect(broadcasts("lobby:match:m-1:chat")).toEqual([]); + }); + + it("rewrites the bell's preview rather than notifying again", async () => { + store(); + + await edit("fixed"); + await flush(); + + expect(notifications.updateChatMessagePreview).toHaveBeenCalledWith( + MESSAGE_ID, + "<b>fixed</b>", + ); + expect(notifications.notifyPlayers).not.toHaveBeenCalled(); + }); + + it("still edits when the preview cannot be rewritten", async () => { + store(); + notifications.updateChatMessagePreview.mockRejectedValue( + new Error("database down"), + ); + + await expect(edit()).resolves.toMatchObject({ edited: true }); + expect(current().message).toBe("fixed"); + }); + + it("never relays an edit to the game server", async () => { + store(); + + await edit(); + await flush(); + + expect(rcon.connect).not.toHaveBeenCalled(); + }); + + it.each([ + ["another player", { steam_id: FRIEND }], + ["an administrator", { steam_id: FRIEND, role: "administrator" }], + ])("refuses %s", async (_, overrides) => { + store(); + role = (overrides as { role?: string }).role ?? "user"; + + await expect(edit("mine now", me(overrides))).resolves.toEqual({ + edited: false, + code: ChatErrorCode.NotAllowed, + }); + expect(current().message).toBe("typo"); + }); + + it.each([ + ["stored before source was recorded", { source: undefined }], + ["relayed from the game", { source: "game" }], + ])("refuses a message %s", async (_, overrides) => { + store(overrides); + + await expect(edit()).resolves.toEqual({ + edited: false, + code: ChatErrorCode.NotAllowed, + }); + expect(current().message).toBe("typo"); + }); + + it("refuses once the window has closed", async () => { + store({ timestamp: ago(WINDOW + 6_000) }); + + await expect(edit()).resolves.toEqual({ + edited: false, + code: ChatErrorCode.WindowClosed, + }); + expect(current().message).toBe("typo"); + }); + + it("allows for a few seconds of clock skew between pods", async () => { + store({ timestamp: ago(WINDOW + 2_000) }); + + await expect(edit()).resolves.toMatchObject({ edited: true }); + }); + + it("keeps a gagged author from editing", async () => { + store(); + gagged = true; + + await expect(edit()).resolves.toEqual({ + edited: false, + code: ChatErrorCode.Gagged, + }); + expect(current().message).toBe("typo"); + }); + + it("refuses an author who can no longer get into the room", async () => { + store({}, "m-2"); + + await expect(edit("fixed", me(), "m-2")).resolves.toEqual({ + edited: false, + code: ChatErrorCode.NotAllowed, + }); + }); + + it("refuses an edit over the limit", async () => { + store(); + + await expect( + edit("a".repeat(ChatService.MAX_MESSAGE_LENGTH + 1)), + ).resolves.toEqual({ edited: false, code: ChatErrorCode.TooLong }); + expect(current().message).toBe("typo"); + }); + + it("answers not_found for a message that is not there", async () => { + await expect(edit()).resolves.toEqual({ + edited: false, + code: ChatErrorCode.NotFound, + }); + }); + + it("answers not_found for an id that could never be a message", async () => { + await expect( + service.editMessage( + ChatLobbyType.Match, + "m-1", + "not-a-uuid", + me(), + "fixed", + ), + ).resolves.toEqual({ edited: false, code: ChatErrorCode.NotFound }); + expect(redis.hget).not.toHaveBeenCalled(); + }); + + it("does not bring back a message deleted between the read and the write", async () => { + store(); + redis.hget.mockImplementationOnce( + async (key: string, field: string) => { + const raw = stored[key][field]; + delete stored[key][field]; + return raw; + }, + ); + + await expect(edit()).resolves.toEqual({ + edited: false, + code: ChatErrorCode.NotFound, + }); + expect(current()).toBeNull(); + expect(broadcasts("lobby:match:m-1:edited")).toEqual([]); + }); + + it("applies an edit on top of one that landed between the read and the write", async () => { + store(); + redis.hget.mockImplementationOnce( + async (key: string, field: string) => { + const raw = stored[key][field]; + stored[key][field] = JSON.stringify({ + ...JSON.parse(raw), + message: "other tab", + edited_at: new Date().toISOString(), + }); + return raw; + }, + ); + + await expect(edit()).resolves.toMatchObject({ edited: true }); + expect(current().message).toBe("fixed"); + }); + + it("lets the author delete their own message, and audits it", async () => { + store(); + + await expect(selfDelete()).resolves.toEqual({ deleted: true }); + + expect(current()).toBeNull(); + expect(audits().at(0)?.bindings).toEqual([ + MESSAGE_ID, + "match", + "m-1", + ME, + "typo", + expect.any(String), + "web", + ME, + ]); + expect(notifications.retractChatMessage).toHaveBeenCalledWith( + MESSAGE_ID, + ); + }); + + it("lets a gagged author delete their own message", async () => { + store(); + gagged = true; + + await expect(selfDelete()).resolves.toEqual({ deleted: true }); + }); + + it("refuses the author's delete once the window has closed", async () => { + store({ timestamp: ago(WINDOW + 6_000) }); + + await expect(selfDelete()).resolves.toEqual({ + deleted: false, + code: ChatErrorCode.WindowClosed, + }); + expect(audits()).toHaveLength(0); + expect(current().message).toBe("typo"); + }); + + it("refuses the author's delete of a line relayed from the game", async () => { + store({ source: "game" }); + + await expect(selfDelete()).resolves.toEqual({ + deleted: false, + code: ChatErrorCode.NotAllowed, + }); + }); + + it("still lets a moderator delete after the window", async () => { + role = "moderator"; + store({ timestamp: ago(WINDOW * 10) }); + + await expect(selfDelete(me({ role: "moderator" }))).resolves.toEqual({ + deleted: true, + }); + }); + }); + + describe("in a direct conversation", () => { + const room = directRoomId(ME, FRIEND); + + const hold = (overrides: Partial = {}) => { + directMessage = { + id: MESSAGE_ID, + roomId: room, + author: ME, + message: "typo", + open: true, + editedAt: null, + ...overrides, + }; + }; + + const editDirect = (user = me()) => + service.editMessage( + ChatLobbyType.Direct, + room, + MESSAGE_ID, + user, + "fixed", + ); + + const deleteDirect = (user = me()) => + service.deleteMessage(ChatLobbyType.Direct, room, MESSAGE_ID, user); + + const statements = (verb: string) => + queries.filter(({ sql }) => + sql.includes(`${verb} public.direct_messages`), + ); + + beforeEach(() => { + redis.hgetall.mockImplementation(async (key: string) => + key === `chat:direct:${room}` + ? { [FRIEND]: JSON.stringify({ user: { steam_id: FRIEND } }) } + : {}, + ); + }); + + it("lets the author edit within the window", async () => { + hold(); + + await expect(editDirect()).resolves.toEqual({ + edited: true, + message: "fixed", + edited_at: "2026-01-01T00:00:00.000Z", + }); + await flush(); + + expect(statements("UPDATE").at(0)?.bindings).toEqual([ + MESSAGE_ID, + room, + ME, + 600, + "fixed", + ]); + expect(broadcasts(`lobby:direct:${room}:edited`)).toEqual([ + { + steamId: FRIEND, + event: `lobby:direct:${room}:edited`, + data: { + id: MESSAGE_ID, + message: "fixed", + edited_at: "2026-01-01T00:00:00.000Z", + }, + }, + ]); + expect(notifications.updateChatMessagePreview).toHaveBeenCalledWith( + MESSAGE_ID, + "fixed", + ); + }); + + it("is not held back by a gag", async () => { + hold(); + gagged = true; + + await expect(editDirect()).resolves.toMatchObject({ edited: true }); + expect( + queries.some(({ sql }) => sql.includes("public.is_gagged")), + ).toBe(false); + }); + + it("refuses the other party's message", async () => { + hold({ author: FRIEND }); + + await expect(editDirect()).resolves.toEqual({ + edited: false, + code: ChatErrorCode.NotAllowed, + }); + expect(statements("UPDATE")).toHaveLength(0); + }); + + it("refuses once the window has closed", async () => { + hold({ open: false }); + + await expect(editDirect()).resolves.toEqual({ + edited: false, + code: ChatErrorCode.WindowClosed, + }); + await expect(deleteDirect()).resolves.toEqual({ + deleted: false, + code: ChatErrorCode.WindowClosed, + }); + expect(statements("UPDATE")).toHaveLength(0); + expect(statements("DELETE FROM")).toHaveLength(0); + }); + + it("refuses once the friendship is gone, before reading anything", async () => { + hold(); + acceptedFriendships = []; + + await expect(editDirect()).resolves.toEqual({ + edited: false, + code: ChatErrorCode.NotAllowed, + }); + await expect(deleteDirect()).resolves.toEqual({ + deleted: false, + code: ChatErrorCode.NotAllowed, + }); + expect(queries).toHaveLength(0); + }); + + it("answers not_found for a message that is not there", async () => { + await expect(editDirect()).resolves.toEqual({ + edited: false, + code: ChatErrorCode.NotFound, + }); + }); + + it("says why when the row changed between the read and the write", async () => { + hold(); + postgres.query.mockImplementationOnce(async (sql, bindings) => { + queries.push({ sql, bindings }); + const row = { author: ME, open: true }; + directMessage = undefined; + return [row]; + }); + + await expect(editDirect()).resolves.toEqual({ + edited: false, + code: ChatErrorCode.NotFound, + }); + expect(broadcasts(`lobby:direct:${room}:edited`)).toEqual([]); + }); + + it("lets the author delete within the window, without an audit", async () => { + hold(); + + await expect(deleteDirect()).resolves.toEqual({ deleted: true }); + await flush(); + + expect(directMessage).toBeUndefined(); + expect(audits()).toHaveLength(0); + expect(broadcasts(`lobby:direct:${room}:deleted`)).toEqual([ + { + steamId: FRIEND, + event: `lobby:direct:${room}:deleted`, + data: { id: MESSAGE_ID }, + }, + ]); + expect(notifications.retractChatMessage).toHaveBeenCalledWith( + MESSAGE_ID, + ); + }); + + it("refuses to delete the other party's message", async () => { + hold({ author: FRIEND }); + + await expect(deleteDirect()).resolves.toEqual({ + deleted: false, + code: ChatErrorCode.NotAllowed, + }); + expect(directMessage).toBeDefined(); + }); + }); + + describe("while its notifications are still being written", () => { + const say = async (type: ChatLobbyType, id: string) => { + redis.hget.mockImplementation(async (key: string, field: string) => + key.startsWith("chat:") + ? JSON.stringify({ user: { steam_id: ME } }) + : (stored[key]?.[field] ?? null), + ); + + const result = await service.sendMessageToChat(type, id, me(), "typo"); + + return result.accepted ? result.messageId : undefined; + }; + + it("gives the rows a room message's edited text", async () => { + tournament.roster = [ME, FRIEND]; + notifications.notifyPlayers.mockImplementationOnce(async () => { + const [key] = redis.hset.mock.calls.at(-1); + const [field, raw] = redis.hset.mock.calls.at(-1).slice(1); + stored[key] = { + [field]: JSON.stringify({ + ...JSON.parse(raw), + message: "fixed", + edited_at: new Date().toISOString(), + }), + }; + }); + + const messageId = await say(ChatLobbyType.Tournament, "t-1"); + await flush(); + await flush(); + + expect(notifications.updateChatMessagePreview).toHaveBeenCalledWith( + messageId, + "fixed", + ); + expect(notifications.retractChatMessage).not.toHaveBeenCalled(); + }); + + it("leaves the rows alone when the message is unchanged", async () => { + tournament.roster = [ME, FRIEND]; + redis.hset.mockImplementationOnce( + async (key: string, field: string, raw: string) => { + stored[key] = { [field]: raw }; + }, + ); + + await say(ChatLobbyType.Tournament, "t-1"); + await flush(); + await flush(); + + expect(notifications.notifyPlayers).toHaveBeenCalled(); + expect(notifications.updateChatMessagePreview).not.toHaveBeenCalled(); + expect(notifications.retractChatMessage).not.toHaveBeenCalled(); + }); + + it("retracts them when a direct message was deleted meanwhile", async () => { + const messageId = await say( + ChatLobbyType.Direct, + directRoomId(ME, FRIEND), + ); + await flush(); + await flush(); + + expect(notifications.retractChatMessage).toHaveBeenCalledWith( + messageId, + ); + }); + + it("gives them a direct message's edited text", async () => { + notifications.notifyPlayers.mockImplementationOnce(async () => { + const insert = queries.find(({ sql }) => + sql.includes("INSERT INTO public.direct_messages"), + ); + directMessage = { + id: insert.bindings[0], + roomId: insert.bindings[1], + author: ME, + message: "fixed", + open: true, + editedAt: new Date(), + }; + }); + + const messageId = await say( + ChatLobbyType.Direct, + directRoomId(ME, FRIEND), + ); + await flush(); + await flush(); + + expect(notifications.updateChatMessagePreview).toHaveBeenCalledWith( + messageId, + "fixed", + ); + expect(notifications.retractChatMessage).not.toHaveBeenCalled(); + }); + }); + }); + describe("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 75d3b01a..6c683b11 100644 --- a/src/chat/chat.service.ts +++ b/src/chat/chat.service.ts @@ -22,6 +22,7 @@ import { ChatErrorCode } from "./enums/ChatErrorCode"; import { ChatMessage, ChatMessageSource } from "./types/ChatMessage"; import { ChatSendResult } from "./types/ChatSendResult"; import { ChatDeleteResult } from "./types/ChatDeleteResult"; +import { ChatEditResult } from "./types/ChatEditResult"; @Injectable() export class ChatService { @@ -52,6 +53,39 @@ export class ChatService { // packet. public static readonly RCON_MESSAGE_MAX_LENGTH = 240; + // How long an author may edit or delete what they sent from the website. + public static readonly SELF_SERVICE_WINDOW_MS = 600_000; + + // A room message carries the clock of whichever pod stored it, and another + // pod may be the one judging the window. + private static readonly SELF_SERVICE_CLOCK_SKEW_MS = 5_000; + + private static readonly EDIT_ATTEMPTS = 3; + + // HSET drops a field's expiry, so the absolute expiry is read first and put + // back: an edit never extends a message's life. Comparing against the value + // that was read and checked keeps an edit from bringing back a message that + // expired or was deleted in the meantime, or from overwriting another edit. + private static readonly EDIT_ROOM_MESSAGE_SCRIPT = ` + if redis.call('HGET', KEYS[1], ARGV[1]) ~= ARGV[2] then + return 0 + end + local expiresAt = redis.call('HPEXPIRETIME', KEYS[1], 'FIELDS', 1, ARGV[1])[1] + redis.call('HSET', KEYS[1], ARGV[1], ARGV[3]) + if expiresAt > 0 then + redis.call('HPEXPIREAT', KEYS[1], expiresAt, 'FIELDS', 1, ARGV[1]) + end + return 1 + `; + + // 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. + private static readonly OWN_RECENT_DIRECT_MESSAGE = `id = $1::uuid + AND room_id = $2 + AND from_steam_id = $3::bigint + AND created_at > now() - make_interval(secs => $4::int)`; + // A drafted free agent is on a roster and gets in that way; withdrawn means // they left the pool. private static readonly TOURNAMENT_CHAT_FREE_AGENT_STATUSES: e_tournament_free_agent_statuses_enum[] = @@ -694,8 +728,12 @@ export class ChatService { messageId: string, user: User, ): Promise { + if (!ChatService.UUID.test(messageId)) { + return { deleted: false, code: ChatErrorCode.NotFound }; + } + if (type === ChatLobbyType.Direct) { - return { deleted: false, code: ChatErrorCode.NotAllowed }; + return await this.deleteDirectMessage(id, messageId, user); } const current = await this.getCurrentUser(user.steam_id); @@ -705,9 +743,7 @@ export class ChatService { } const messageKey = `chat_${type}_${id}`; - const raw = ChatService.UUID.test(messageId) - ? await this.redis.hget(messageKey, messageId) - : null; + const raw = await this.redis.hget(messageKey, messageId); if (!raw) { return { deleted: false, code: ChatErrorCode.NotFound }; @@ -715,39 +751,338 @@ export class ChatService { const message = JSON.parse(raw) as ChatMessage; - if (!(await this.canDelete(message, type, id, current))) { - return { deleted: false, code: ChatErrorCode.NotAllowed }; + const refusal = await this.deleteRefusal(message, type, id, current); + + if (refusal) { + return { deleted: false, code: refusal }; } // Audited before it is removed, so no failure part way through can take a // message down without its evidence. A retry finds the row and carries on. + // An author removing their own message is audited the same way, or posting + // abuse and deleting it would leave nothing behind. await this.recordDeletion(type, id, messageId, message, current); await this.redis.hdel(messageKey, messageId); void this.to(type, id, "deleted", { id: messageId }); - await this.notifications.retractChatMessage(messageId).catch((error) => { - this.logger.warn( - `unable to retract notifications for ${type}:${id} message ${messageId}`, - error, - ); - }); + await this.retractNotifications(type, id, messageId); return { deleted: true }; } - private async canDelete( + private async deleteRefusal( message: ChatMessage, type: ChatLobbyType, id: string, user: User, - ): Promise { + ): Promise { if (!isRoleAbove(user.role, "moderator")) { - return false; + const refusal = ChatService.selfServiceRefusal(message, user); + + if (refusal) { + return refusal; + } } - return await this.canAccessLobby(type, id, user); + if (!(await this.canAccessLobby(type, id, user))) { + return ChatErrorCode.NotAllowed; + } + + return null; + } + + // Only what the author typed on the website is theirs to change: a line + // relayed from the game, or one stored before `source` was recorded, is not. + // No role gets past this -- nobody edits another player's words. + private static selfServiceRefusal( + message: ChatMessage, + user: User, + ): ChatErrorCode | null { + if ( + message.source !== "web" || + ChatService.authorSteamId(message) !== String(user.steam_id) + ) { + return ChatErrorCode.NotAllowed; + } + + const sentAt = new Date(message.timestamp).getTime(); + + if ( + Number.isNaN(sentAt) || + Date.now() - sentAt > + ChatService.SELF_SERVICE_WINDOW_MS + + ChatService.SELF_SERVICE_CLOCK_SKEW_MS + ) { + return ChatErrorCode.WindowClosed; + } + + return null; + } + + public async editMessage( + type: ChatLobbyType, + id: string, + messageId: string, + user: User, + raw: unknown, + ): Promise { + const parsed = ChatService.messageText(raw); + + if ("error" in parsed) { + return { edited: false, code: parsed.error }; + } + + if (!ChatService.UUID.test(messageId)) { + return { edited: false, code: ChatErrorCode.NotFound }; + } + + if (type === ChatLobbyType.Direct) { + return await this.editDirectMessage(id, messageId, user, parsed.text); + } + + const current = await this.getCurrentUser(user.steam_id); + + if (!current) { + return { edited: false, code: ChatErrorCode.NotAllowed }; + } + + return await this.editRoomMessage( + type, + id, + messageId, + current, + parsed.text, + ); + } + + private async editRoomMessage( + type: ChatLobbyType, + id: string, + messageId: string, + user: User, + text: string, + ): Promise { + const messageKey = `chat_${type}_${id}`; + let admitted = false; + + for (let attempt = 0; attempt < ChatService.EDIT_ATTEMPTS; attempt++) { + const raw = await this.redis.hget(messageKey, messageId); + + if (!raw) { + return { edited: false, code: ChatErrorCode.NotFound }; + } + + const message = JSON.parse(raw) as ChatMessage; + const refusal = ChatService.selfServiceRefusal(message, user); + + if (refusal) { + return { edited: false, code: refusal }; + } + + if (!admitted) { + if (!(await this.canAccessLobby(type, id, user))) { + return { edited: false, code: ChatErrorCode.NotAllowed }; + } + + if (await this.isGagged(user.steam_id)) { + return { edited: false, code: ChatErrorCode.Gagged }; + } + + admitted = true; + } + + const editedAt = new Date().toISOString(); + + const swapped = await this.redis.eval( + ChatService.EDIT_ROOM_MESSAGE_SCRIPT, + 1, + messageKey, + messageId, + raw, + JSON.stringify({ ...message, message: text, edited_at: editedAt }), + ); + + if (swapped === 1) { + return await this.announceEdit(type, id, messageId, text, editedAt); + } + } + + this.logger.warn( + `gave up editing ${type}:${id} message ${messageId}, it kept changing`, + ); + + return { edited: false, code: ChatErrorCode.Invalid }; + } + + private async editDirectMessage( + roomId: string, + messageId: string, + user: User, + text: string, + ): Promise { + if (!(await this.canAccessLobby(ChatLobbyType.Direct, roomId, user))) { + return { edited: false, code: ChatErrorCode.NotAllowed }; + } + + const refusal = await this.directMessageRefusal(roomId, messageId, user); + + if (refusal) { + return { edited: false, code: refusal }; + } + + const [row] = await this.postgres.query< + Array<{ message: string; edited_at: Date }> + >( + `UPDATE public.direct_messages + SET message = $5, edited_at = now() + WHERE ${ChatService.OWN_RECENT_DIRECT_MESSAGE} + RETURNING message, edited_at`, + [...ChatService.ownRecentDirectMessage(roomId, messageId, user), text], + ); + + if (!row) { + return { + edited: false, + code: + (await this.directMessageRefusal(roomId, messageId, user)) ?? + ChatErrorCode.NotFound, + }; + } + + return await this.announceEdit( + ChatLobbyType.Direct, + roomId, + messageId, + row.message, + new Date(row.edited_at).toISOString(), + ); + } + + private async deleteDirectMessage( + roomId: string, + messageId: string, + user: User, + ): Promise { + if (!(await this.canAccessLobby(ChatLobbyType.Direct, roomId, user))) { + return { deleted: false, code: ChatErrorCode.NotAllowed }; + } + + const refusal = await this.directMessageRefusal(roomId, messageId, user); + + if (refusal) { + return { deleted: false, code: refusal }; + } + + const [row] = await this.postgres.query>( + `DELETE FROM public.direct_messages + WHERE ${ChatService.OWN_RECENT_DIRECT_MESSAGE} + RETURNING id::text AS id`, + ChatService.ownRecentDirectMessage(roomId, messageId, user), + ); + + if (!row) { + return { + deleted: false, + code: + (await this.directMessageRefusal(roomId, messageId, user)) ?? + ChatErrorCode.NotFound, + }; + } + + void this.to(ChatLobbyType.Direct, roomId, "deleted", { id: messageId }); + + await this.retractNotifications(ChatLobbyType.Direct, roomId, messageId); + + return { deleted: true }; + } + + private static ownRecentDirectMessage( + roomId: string, + messageId: string, + user: User, + ) { + return [ + messageId, + roomId, + String(user.steam_id), + ChatService.SELF_SERVICE_WINDOW_MS / 1000, + ]; + } + + // Read apart from the statement that acts on it, so a refusal can say why. + private async directMessageRefusal( + roomId: string, + messageId: string, + user: User, + ): Promise { + const [row] = await this.postgres.query< + Array<{ author: string; open: boolean }> + >( + `SELECT from_steam_id::text AS author, + created_at > now() - make_interval(secs => $3::int) AS open + FROM public.direct_messages + WHERE id = $1::uuid AND room_id = $2`, + [messageId, roomId, ChatService.SELF_SERVICE_WINDOW_MS / 1000], + ); + + if (!row) { + return ChatErrorCode.NotFound; + } + + if (row.author !== String(user.steam_id)) { + return ChatErrorCode.NotAllowed; + } + + if (!row.open) { + return ChatErrorCode.WindowClosed; + } + + return null; + } + + // Never relayed to the game server, and no new notification: the bell rows + // already announcing the message just show what it says now. + private async announceEdit( + type: ChatLobbyType, + id: string, + messageId: string, + text: string, + editedAt: string, + ): Promise { + void this.to(type, id, "edited", { + id: messageId, + message: text, + edited_at: editedAt, + }); + + await this.notifications + .updateChatMessagePreview( + messageId, + ChatService.notificationPreview(text), + ) + .catch((error) => { + this.logger.warn( + `unable to update notifications for ${type}:${id} message ${messageId}`, + error, + ); + }); + + return { edited: true, message: text, edited_at: editedAt }; + } + + private async retractNotifications( + type: ChatLobbyType, + id: string, + messageId: string, + ) { + await this.notifications.retractChatMessage(messageId).catch((error) => { + this.logger.warn( + `unable to retract notifications for ${type}:${id} message ${messageId}`, + error, + ); + }); } private async recordDeletion( @@ -839,9 +1174,7 @@ export class ChatService { await this.notifications.notifyPlayers(notificationType, { title: senderName, - message: NotificationsService.escapeHtml( - message.length > 140 ? `${message.slice(0, 140)}…` : message, - ), + message: ChatService.notificationPreview(message), role: "user", entity_id: entityId, steamIds: targets, @@ -864,14 +1197,67 @@ export class ChatService { targets, ); - // A delete that landed while the rows above were being written retracted - // nothing. Its audit row is committed before it retracts, and unlike the - // redis field it is not gone just because the message expired or moved. - if (type !== ChatLobbyType.Direct && (await this.wasDeleted(messageId))) { + await this.catchUpNotifications(type, id, messageId); + } + + // A delete or an edit that landed while the rows were being written had no + // rows to act on yet, so whatever became of the message is applied now. + // + // A room's audit row is what says it was deleted: unlike the redis field, it + // is not gone just because the message expired or moved. A direct message + // has no audit, but one written moments ago cannot have been pruned yet, so + // a missing row was deleted. + private async catchUpNotifications( + type: ChatLobbyType, + id: string, + messageId: string, + ) { + if (type === ChatLobbyType.Direct) { + const [row] = await this.postgres.query< + Array<{ message: string; edited_at: Date | null }> + >( + `SELECT message, edited_at FROM public.direct_messages + WHERE id = $1::uuid`, + [messageId], + ); + + if (!row) { + await this.notifications.retractChatMessage(messageId); + return; + } + + if (row.edited_at) { + await this.notifications.updateChatMessagePreview( + messageId, + ChatService.notificationPreview(row.message), + ); + } + + return; + } + + if (await this.wasDeleted(messageId)) { await this.notifications.retractChatMessage(messageId); + return; + } + + const raw = await this.redis.hget(`chat_${type}_${id}`, messageId); + const message = raw ? (JSON.parse(raw) as ChatMessage) : null; + + if (message?.edited_at) { + await this.notifications.updateChatMessagePreview( + messageId, + ChatService.notificationPreview(message.message), + ); } } + private static notificationPreview(message: string) { + return NotificationsService.escapeHtml( + message.length > 140 ? `${message.slice(0, 140)}…` : message, + ); + } + private async wasDeleted(messageId: string): Promise { const [row] = await this.postgres.query>( `SELECT EXISTS ( @@ -1395,6 +1781,7 @@ export class ChatService { id: string; message: string; created_at: Date; + edited_at: Date | null; steam_id: string; name: string; role: e_player_roles_enum; @@ -1402,7 +1789,7 @@ export class ChatService { profile_url: string | null; }> >( - `SELECT dm.id::text AS id, dm.message, dm.created_at, + `SELECT dm.id::text AS id, dm.message, dm.created_at, dm.edited_at, 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 @@ -1420,6 +1807,9 @@ export class ChatService { timestamp: new Date(row.created_at).toISOString(), // Nothing relays from the game into a DM. source: "web", + ...(row.edited_at + ? { edited_at: new Date(row.edited_at).toISOString() } + : {}), from: { role: row.role, name: row.name, @@ -1580,7 +1970,14 @@ export class ChatService { public async to( type: ChatLobbyType, id: string, - event: "chat" | "deleted" | "list" | "messages" | "joined" | "left", + event: + | "chat" + | "edited" + | "deleted" + | "list" + | "messages" + | "joined" + | "left", data: Record, ) { const users = await this.getAllUsersInLobby(type, id); diff --git a/src/chat/enums/ChatErrorCode.ts b/src/chat/enums/ChatErrorCode.ts index 35624e17..00aae33a 100644 --- a/src/chat/enums/ChatErrorCode.ts +++ b/src/chat/enums/ChatErrorCode.ts @@ -6,4 +6,5 @@ export enum ChatErrorCode { Invalid = "invalid", Gagged = "gagged", NotFound = "not_found", + WindowClosed = "window_closed", } diff --git a/src/chat/types/ChatAction.ts b/src/chat/types/ChatAction.ts index 84411874..029223aa 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"; +export type ChatAction = "send" | "delete" | "edit"; diff --git a/src/chat/types/ChatEditResult.ts b/src/chat/types/ChatEditResult.ts new file mode 100644 index 00000000..4faa3d53 --- /dev/null +++ b/src/chat/types/ChatEditResult.ts @@ -0,0 +1,5 @@ +import { ChatErrorCode } from "../enums/ChatErrorCode"; + +export type ChatEditResult = + | { edited: true; message: string; edited_at: string } + | { edited: false; code: ChatErrorCode }; diff --git a/src/chat/types/ChatMessage.ts b/src/chat/types/ChatMessage.ts index 70b4ce46..0ba56e32 100644 --- a/src/chat/types/ChatMessage.ts +++ b/src/chat/types/ChatMessage.ts @@ -8,6 +8,8 @@ export interface ChatMessage { timestamp: string; // Absent on messages stored before it was recorded. source?: ChatMessageSource; + // ISO 8601, present once the author has edited the message. + edited_at?: string; from: { role: e_player_roles_enum; name: string; diff --git a/src/notifications/notifications.service.ts b/src/notifications/notifications.service.ts index f54f0466..14fcc1be 100644 --- a/src/notifications/notifications.service.ts +++ b/src/notifications/notifications.service.ts @@ -789,6 +789,20 @@ export class NotificationsService { ); } + // `deleted_at IS NULL` is what keeps an edit racing a delete from writing + // text back into a row retractChatMessage has already blanked. + async updateChatMessagePreview(messageId: string, preview: string) { + await this.postgres.query( + `UPDATE public.notifications + SET message = $2 + WHERE data->>'messageId' = $1 + AND type IN ('ChatMessage', 'MatchChatMessage') + AND is_read = false + AND deleted_at IS NULL`, + [messageId, preview], + ); + } + // Opening a conversation should clear its badge everywhere, not just in the // tab that was open. async markConversationRead( diff --git a/test/chat-direct-messages.spec.ts b/test/chat-direct-messages.spec.ts index 94869e44..ad53c8df 100644 --- a/test/chat-direct-messages.spec.ts +++ b/test/chat-direct-messages.spec.ts @@ -1,8 +1,12 @@ +import { readFileSync } from "fs"; +import { join } from "path"; import { PostgresService } from "./../src/postgres/postgres.service"; import { Fixtures } from "./utils/fixtures"; import { bootMigratedDb, SqlTestDb } from "./utils/sql-test-db"; import { ChatService } from "./../src/chat/chat.service"; +import { ChatErrorCode } from "./../src/chat/enums/ChatErrorCode"; import { ChatLobbyType } from "./../src/chat/enums/ChatLobbyTypes"; +import { NotificationsService } from "./../src/notifications/notifications.service"; import { PruneDirectMessages } from "./../src/chat/jobs/PruneDirectMessages"; import { directRoomId } from "./../src/chat/utilities/directRoomId"; @@ -14,6 +18,7 @@ describe("direct messages (SQL-driven)", () => { let postgres: PostgresService; let fx: Fixtures; let chat: ChatService; + let bell: NotificationsService; const logger = { log: jest.fn(), warn: jest.fn(), error: jest.fn() }; @@ -36,6 +41,17 @@ describe("direct messages (SQL-driven)", () => { postgres = db.postgres; fx = new Fixtures(postgres, 76561199400000000n); + bell = new NotificationsService( + {} as any, + postgres, + logger as any, + { get: () => ({ webDomain: "https://example.com" }) } as any, + {} as any, + {} as any, + {} as any, + {} as any, + ); + chat = new ChatService( logger as any, {} as any, @@ -49,7 +65,15 @@ describe("direct messages (SQL-driven)", () => { } as any, postgres, { getConnection: () => redis } as any, - { notifyPlayers: jest.fn(), markConversationRead: jest.fn() } as any, + { + notifyPlayers: jest.fn(), + markConversationRead: jest.fn(), + collapseOlderUnread: jest.fn(), + retractChatMessage: (messageId: string) => + bell.retractChatMessage(messageId), + updateChatMessagePreview: (messageId: string, preview: string) => + bell.updateChatMessagePreview(messageId, preview), + } as any, ); }, 600_000); @@ -391,4 +415,242 @@ describe("direct messages (SQL-driven)", () => { expect(row.count).toBe("1"); }); }); + + describe("editing and deleting your own messages", () => { + const sent = async (roomId: string, from: string, message = "typo") => { + const result = await say(roomId, from, message); + return result.accepted ? result.messageId : ""; + }; + + const stored = async (id: string) => + ( + await postgres.query< + Array<{ message: string; created_at: Date; edited_at: Date | null }> + >( + `SELECT message, created_at, edited_at FROM direct_messages + WHERE id = $1::uuid`, + [id], + ) + ).at(0); + + const age = (id: string, minutes: number) => + postgres.query( + `UPDATE direct_messages + SET created_at = created_at - make_interval(mins => $2::int) + WHERE id = $1::uuid`, + [id, minutes], + ); + + const bellRow = async (steamId: string, messageId: string) => { + const [row] = await postgres.query>( + `INSERT INTO notifications + (type, title, message, role, steam_id, entity_id, data) + VALUES ('ChatMessage', 'Someone', 'typo', 'user', $1::bigint, + 'direct:' || $1, jsonb_build_object('messageId', $2::text)) + RETURNING id::text AS id`, + [steamId, messageId], + ); + return row.id; + }; + + const notification = async (id: string) => + ( + await postgres.query< + Array<{ message: string; deleted_at: Date | null }> + >(`SELECT message, deleted_at FROM notifications WHERE id = $1::uuid`, [ + id, + ]) + ).at(0); + + const as = (steamId: string) => ({ steam_id: steamId }) as any; + + it("stamps edited_at and leaves created_at alone", async () => { + const me = await fx.player(); + const friend = await fx.player(); + const room = directRoomId(me, friend); + const id = await sent(room, me); + const before = await stored(id); + + const result = await chat.editMessage( + ChatLobbyType.Direct, + room, + id, + as(me), + " fixed ", + ); + const after = await stored(id); + + expect(before.edited_at).toBeNull(); + expect(after.message).toBe("fixed"); + expect(after.created_at).toEqual(before.created_at); + expect(result).toEqual({ + edited: true, + message: "fixed", + edited_at: after.edited_at.toISOString(), + }); + }); + + it("shows the edit in history after a reload", async () => { + const me = await fx.player(); + const friend = await fx.player(); + const room = directRoomId(me, friend); + const edited = await sent(room, me); + await sent(room, friend, "untouched"); + + await chat.editMessage( + ChatLobbyType.Direct, + room, + edited, + as(me), + "fixed", + ); + + const [first, second] = await chat["getDirectMessages"](room); + + expect(first).toMatchObject({ + id: edited, + message: "fixed", + edited_at: (await stored(edited)).edited_at.toISOString(), + }); + expect(second).not.toHaveProperty("edited_at"); + }); + + it("refuses an edit or a delete once the message is older than the window", async () => { + const me = await fx.player(); + const friend = await fx.player(); + const room = directRoomId(me, friend); + const id = await sent(room, me); + await age(id, 11); + + await expect( + chat.editMessage(ChatLobbyType.Direct, room, id, as(me), "fixed"), + ).resolves.toEqual({ edited: false, code: ChatErrorCode.WindowClosed }); + await expect( + chat.deleteMessage(ChatLobbyType.Direct, room, id, as(me)), + ).resolves.toEqual({ deleted: false, code: ChatErrorCode.WindowClosed }); + + expect((await stored(id))?.message).toBe("typo"); + }); + + it("refuses the other party", async () => { + const me = await fx.player(); + const friend = await fx.player(); + const room = directRoomId(me, friend); + const id = await sent(room, me); + + await expect( + chat.editMessage(ChatLobbyType.Direct, room, id, as(friend), "mine"), + ).resolves.toEqual({ edited: false, code: ChatErrorCode.NotAllowed }); + await expect( + chat.deleteMessage(ChatLobbyType.Direct, room, id, as(friend)), + ).resolves.toEqual({ deleted: false, code: ChatErrorCode.NotAllowed }); + + expect((await stored(id))?.message).toBe("typo"); + }); + + it("answers not_found for a message from another conversation", async () => { + const me = await fx.player(); + const friend = await fx.player(); + const other = await fx.player(); + const id = await sent(directRoomId(me, other), me); + + await expect( + chat.editMessage( + ChatLobbyType.Direct, + directRoomId(me, friend), + id, + as(me), + "fixed", + ), + ).resolves.toEqual({ edited: false, code: ChatErrorCode.NotFound }); + }); + + it("shows the edit on the recipient's unread bell row", async () => { + const me = await fx.player(); + const friend = await fx.player(); + const room = directRoomId(me, friend); + const id = await sent(room, me); + const row = await bellRow(friend, id); + + await chat.editMessage(ChatLobbyType.Direct, room, id, as(me), "fixed"); + + expect(await notification(row)).toEqual({ + message: "fixed", + deleted_at: null, + }); + }); + + it("deletes within the window and retracts the recipient's bell row", async () => { + const me = await fx.player(); + const friend = await fx.player(); + const room = directRoomId(me, friend); + const id = await sent(room, me); + const row = await bellRow(friend, id); + + await expect( + chat.deleteMessage(ChatLobbyType.Direct, room, id, as(me)), + ).resolves.toEqual({ deleted: true }); + + expect(await stored(id)).toBeUndefined(); + expect(await notification(row)).toEqual({ + message: "", + deleted_at: expect.any(Date), + }); + + const [{ count }] = await postgres.query>( + `SELECT count(*)::text AS count FROM chat_message_deletions`, + ); + expect(count).toBe("0"); + }); + + it("answers not_found once it is gone", async () => { + const me = await fx.player(); + const friend = await fx.player(); + const room = directRoomId(me, friend); + const id = await sent(room, me); + + await chat.deleteMessage(ChatLobbyType.Direct, room, id, as(me)); + + await expect( + chat.deleteMessage(ChatLobbyType.Direct, room, id, as(me)), + ).resolves.toEqual({ deleted: false, code: ChatErrorCode.NotFound }); + await expect( + chat.editMessage(ChatLobbyType.Direct, room, id, as(me), "fixed"), + ).resolves.toEqual({ edited: false, code: ChatErrorCode.NotFound }); + }); + }); + + describe("the edited_at migration", () => { + const migration = (file: string) => + readFileSync( + join( + __dirname, + "../hasura/migrations/default/1888000000200_direct_messages_edited_at", + file, + ), + "utf8", + ); + + const hasColumn = async () => + ( + await postgres.query>( + `SELECT column_name FROM information_schema.columns + WHERE table_schema = 'public' + AND table_name = 'direct_messages' + AND column_name = 'edited_at'`, + ) + ).length === 1; + + it("re-applies cleanly and rolls back", async () => { + await postgres.query(migration("up.sql")); + await postgres.query(migration("up.sql")); + expect(await hasColumn()).toBe(true); + + await postgres.query(migration("down.sql")); + expect(await hasColumn()).toBe(false); + + await postgres.query(migration("up.sql")); + expect(await hasColumn()).toBe(true); + }); + }); }); diff --git a/test/chat-redis-actions.spec.ts b/test/chat-redis-actions.spec.ts new file mode 100644 index 00000000..1f121a04 --- /dev/null +++ b/test/chat-redis-actions.spec.ts @@ -0,0 +1,469 @@ +import { randomUUID } from "crypto"; +import IORedis, { Redis } from "ioredis"; +import { GenericContainer, StartedTestContainer } from "testcontainers"; +import { PostgresService } from "./../src/postgres/postgres.service"; +import { Fixtures } from "./utils/fixtures"; +import { bootMigratedDb, SqlTestDb } from "./utils/sql-test-db"; +import { ChatService } from "./../src/chat/chat.service"; +import { ChatErrorCode } from "./../src/chat/enums/ChatErrorCode"; +import { ChatLobbyType } from "./../src/chat/enums/ChatLobbyTypes"; +import { NotificationsService } from "./../src/notifications/notifications.service"; +import { NotificationPreferencesService } from "./../src/notifications/preferences/notification-preferences.service"; + +// The edit is one compare-and-set script against real hash-field expiry, which +// no fake reproduces: HSET dropping a field's TTL is the whole reason it exists. +describe("chat edits and self deletes (SQL-driven)", () => { + let db: SqlTestDb; + let postgres: PostgresService; + let fx: Fixtures; + let container: StartedTestContainer; + let redis: Redis; + let chat: ChatService; + let notifications: NotificationsService; + + const logger = { log: jest.fn(), warn: jest.fn(), error: jest.fn() }; + + const hasura = () => ({ + query: jest.fn(async (query: any) => { + if (query.players_by_pk) { + const [player] = await postgres.query< + Array<{ steam_id: string; name: string; role: string }> + >( + `SELECT steam_id::text AS steam_id, name, role::text AS role + FROM players WHERE steam_id = $1::bigint`, + [query.players_by_pk.__args.steam_id], + ); + return { players_by_pk: player ?? null }; + } + + // Who gets into which room is chat.service.spec's subject. + if (query.matches_by_pk) { + return { + matches_by_pk: { + is_coach: false, + is_organizer: true, + is_in_lineup: false, + }, + }; + } + + if (query.draft_games) { + return { draft_games: [{ id: query.draft_games.__args.where.id._eq }] }; + } + + return {}; + }), + }); + + beforeAll(async () => { + container = await new GenericContainer("redis:8.8-alpine") + .withExposedPorts(6379) + .start(); + redis = new IORedis({ + host: container.getHost(), + port: container.getMappedPort(6379), + }); + + db = await bootMigratedDb("ChatRedisActionsTest"); + postgres = db.postgres; + fx = new Fixtures(postgres, 76561199610000000n); + + notifications = new NotificationsService( + hasura() as any, + postgres, + logger as any, + { get: () => ({ webDomain: "https://example.com" }) } as any, + new NotificationPreferencesService(postgres), + { add: jest.fn() } as any, + { add: jest.fn() } as any, + { add: jest.fn() } as any, + ); + + chat = new ChatService( + logger as any, + {} as any, + hasura() as any, + postgres, + { getConnection: () => redis } as any, + notifications, + ); + }, 600_000); + + afterAll(async () => { + redis?.disconnect(); + await container?.stop(); + await db?.stop(); + }); + + beforeEach(async () => { + jest.restoreAllMocks(); + jest.clearAllMocks(); + await redis.flushall(); + await postgres.query("DELETE FROM chat_message_deletions"); + await postgres.query("DELETE FROM player_sanctions"); + await postgres.query("DELETE FROM notifications"); + await postgres.query("DELETE FROM players"); + }); + + const author = async () => { + const steamId = await fx.player("Author"); + return { steam_id: steamId, name: "Author", role: "user" } as any; + }; + + const key = (matchId: string) => `chat_match_${matchId}`; + + const written = (id: string, from: any) => + JSON.stringify({ + id, + message: "typo", + timestamp: new Date().toISOString(), + source: "web", + from: { + role: "user", + name: "Author", + steam_id: from.steam_id, + avatar_url: null, + profile_url: "https://steamcommunity.com/profiles/1/", + }, + }); + + const place = async (matchId: string, from: any) => { + const id = randomUUID(); + await redis.hset(key(matchId), id, written(id, from)); + return id; + }; + + const expiresAt = async (matchId: string, id: string) => + ( + (await redis.call( + "HPEXPIRETIME", + key(matchId), + "FIELDS", + 1, + id, + )) as number[] + )[0]; + + const edit = (matchId: string, id: string, user: any, text = "fixed") => + chat.editMessage(ChatLobbyType.Match, matchId, id, user, text); + + // Holds the edit between its read and its write, which is exactly where a + // delete or an expiry has to be shown not to come back. + const pauseFirstRead = (during: () => Promise) => { + const read = redis.hget.bind(redis); + + jest.spyOn(redis, "hget").mockImplementationOnce((async ( + hashKey: string, + field: string, + ) => { + const value = await read(hashKey, field); + await during(); + return value; + }) as any); + }; + + describe("expiry", () => { + it("keeps the expiry a sent message was given", async () => { + const user = await author(); + const matchId = randomUUID(); + + const sent = await chat.sendMessageToChat( + ChatLobbyType.Match, + matchId, + user, + "typo", + true, + ); + const id = sent.accepted ? sent.messageId : ""; + const before = await expiresAt(matchId, id); + + await new Promise((resolve) => setTimeout(resolve, 20)); + + await expect(edit(matchId, id, user)).resolves.toMatchObject({ + edited: true, + }); + + expect(before).toBeGreaterThan(Date.now()); + expect(await expiresAt(matchId, id)).toBe(before); + }); + + it("keeps an expiry to the millisecond", async () => { + const user = await author(); + const matchId = randomUUID(); + const id = await place(matchId, user); + const at = Date.now() + 123_457; + await redis.call("HPEXPIREAT", key(matchId), at, "FIELDS", 1, id); + + await edit(matchId, id, user); + + expect(await expiresAt(matchId, id)).toBe(at); + }); + + it("gives a message with no expiry none", async () => { + const user = await author(); + const matchId = randomUUID(); + const id = await place(matchId, user); + + await edit(matchId, id, user); + + expect(await expiresAt(matchId, id)).toBe(-1); + expect(JSON.parse(await redis.hget(key(matchId), id)).message).toBe( + "fixed", + ); + }); + + it("never brings back a message that expired between the read and the write", async () => { + const user = await author(); + const matchId = randomUUID(); + const id = await place(matchId, user); + await redis.call("HPEXPIRE", key(matchId), 150, "FIELDS", 1, id); + + pauseFirstRead(() => new Promise((resolve) => setTimeout(resolve, 300))); + + await expect(edit(matchId, id, user)).resolves.toEqual({ + edited: false, + code: ChatErrorCode.NotFound, + }); + + expect(await redis.hexists(key(matchId), id)).toBe(0); + expect(await expiresAt(matchId, id)).toBe(-2); + }); + + it("never brings back a message deleted between the read and the write", async () => { + const user = await author(); + const matchId = randomUUID(); + const id = await place(matchId, user); + + pauseFirstRead(() => redis.hdel(key(matchId), id)); + + await expect(edit(matchId, id, user)).resolves.toEqual({ + edited: false, + code: ChatErrorCode.NotFound, + }); + + expect(await redis.hexists(key(matchId), id)).toBe(0); + }); + + it("writes nothing for a value that has moved on", async () => { + const user = await author(); + const matchId = randomUUID(); + const id = await place(matchId, user); + const current = await redis.hget(key(matchId), id); + const at = Date.now() + 60_000; + await redis.call("HPEXPIREAT", key(matchId), at, "FIELDS", 1, id); + + const script = (ChatService as any).EDIT_ROOM_MESSAGE_SCRIPT; + + await expect( + redis.eval(script, 1, key(matchId), id, "stale", "replacement"), + ).resolves.toBe(0); + await expect( + redis.eval(script, 1, key(matchId), randomUUID(), "", "replacement"), + ).resolves.toBe(0); + + expect(await redis.hget(key(matchId), id)).toBe(current); + expect(await expiresAt(matchId, id)).toBe(at); + expect(await redis.hlen(key(matchId))).toBe(1); + }); + }); + + it("never resurrects a message when an edit races its deletion", async () => { + const user = await author(); + + for (let round = 0; round < 25; round++) { + const matchId = randomUUID(); + const id = await place(matchId, user); + + const [edited, deleted] = await Promise.all([ + edit(matchId, id, user), + chat.deleteMessage(ChatLobbyType.Match, matchId, id, user), + ]); + + expect(deleted).toEqual({ deleted: true }); + expect([true, false]).toContain(edited.edited); + expect(await redis.hexists(key(matchId), id)).toBe(0); + } + + const [{ count }] = await postgres.query>( + `SELECT count(*)::text AS count FROM chat_message_deletions + WHERE deleted_by_steam_id = $1::bigint`, + [user.steam_id], + ); + + expect(count).toBe("25"); + }); + + it("stores the edit exactly as JSON.stringify writes it", async () => { + const user = await author(); + const matchId = randomUUID(); + const id = await place(matchId, user); + const original = JSON.parse(await redis.hget(key(matchId), id)); + const text = `héllo — 你好 🎉 a/b "quoted" \\ back\nline`; + + const result = await edit(matchId, id, user, text); + const raw = await redis.hget(key(matchId), id); + + expect(result.edited).toBe(true); + expect(raw).toBe( + JSON.stringify({ + ...original, + message: text, + edited_at: result.edited ? result.edited_at : undefined, + }), + ); + expect(raw).not.toContain("\\/"); + expect(JSON.parse(raw).from.avatar_url).toBeNull(); + expect(JSON.parse(raw).message).toBe(text); + }); + + it("carries the edit through history and a draft moving into its match", async () => { + const user = await author(); + const draftId = randomUUID(); + const matchId = randomUUID(); + const id = randomUUID(); + await redis.hset(`chat_draft_${draftId}`, id, written(id, user)); + + const result = await chat.editMessage( + ChatLobbyType.Draft, + draftId, + id, + user, + "fixed", + ); + + await chat.migrateLobbyMessages( + ChatLobbyType.Draft, + draftId, + ChatLobbyType.Match, + matchId, + ); + + expect(await chat["getMessages"](ChatLobbyType.Match, matchId)).toEqual([ + expect.objectContaining({ + id, + message: "fixed", + edited_at: result.edited ? result.edited_at : undefined, + }), + ]); + }); + + it("audits an author deleting their own message", async () => { + const user = await author(); + const matchId = randomUUID(); + const id = await place(matchId, user); + + await expect( + chat.deleteMessage(ChatLobbyType.Match, matchId, id, user), + ).resolves.toEqual({ deleted: true }); + + const rows = await postgres.query< + Array<{ author: string; deleted_by: string; message: string }> + >( + `SELECT author_steam_id::text AS author, + deleted_by_steam_id::text AS deleted_by, message + FROM chat_message_deletions WHERE message_id = $1::uuid`, + [id], + ); + + expect(rows).toEqual([ + { author: user.steam_id, deleted_by: user.steam_id, message: "typo" }, + ]); + }); + + describe("the bell's preview", () => { + const row = async ( + steamId: string, + messageId: string, + overrides: { + is_read?: boolean; + deleted?: boolean; + message?: string; + } = {}, + ) => { + const [inserted] = await postgres.query>( + `INSERT INTO notifications + (type, title, message, role, steam_id, entity_id, data, + is_read, deleted_at) + VALUES ('MatchChatMessage', 'Author', $3, 'user', $1::bigint, + 'match:m-1', + jsonb_build_object('messageId', $2::text), + $4, CASE WHEN $5 THEN now() END) + RETURNING id::text AS id`, + [ + steamId, + messageId, + overrides.message ?? "typo", + overrides.is_read ?? false, + overrides.deleted ?? false, + ], + ); + return inserted.id; + }; + + const text = async (id: string) => + ( + await postgres.query>( + `SELECT message FROM notifications WHERE id = $1::uuid`, + [id], + ) + )[0].message; + + it("shows the edited text on the recipient's unread row", async () => { + const user = await author(); + const reader = await fx.player("Reader"); + const matchId = randomUUID(); + const id = await place(matchId, user); + const unread = await row(reader, id); + + await edit(matchId, id, user, "fixed"); + + expect(await text(unread)).toBe("<b>fixed</b>"); + }); + + it("rewrites only unread, live rows for that message", async () => { + const reader = await fx.player("Reader"); + const messageId = randomUUID(); + const unread = await row(reader, messageId); + const read = await row(reader, messageId, { is_read: true }); + const collapsed = await row(reader, messageId, { deleted: true }); + const other = await row(reader, randomUUID()); + + await notifications.updateChatMessagePreview(messageId, "fixed"); + + expect(await text(unread)).toBe("fixed"); + expect(await text(read)).toBe("typo"); + expect(await text(collapsed)).toBe("typo"); + expect(await text(other)).toBe("typo"); + }); + + it("leaves a retracted row blank when the edit lands after the delete", async () => { + const reader = await fx.player("Reader"); + const messageId = randomUUID(); + const id = await row(reader, messageId); + + await notifications.retractChatMessage(messageId); + await notifications.updateChatMessagePreview(messageId, "fixed"); + + expect(await text(id)).toBe(""); + }); + + it("finds the message's rows through an index", async () => { + const plan = await postgres.transaction(async (client) => { + await client.query("SET LOCAL enable_seqscan = off"); + + const { rows } = await client.query( + `EXPLAIN UPDATE notifications SET message = 'x' + WHERE data->>'messageId' = $1 + AND type IN ('ChatMessage', 'MatchChatMessage') + AND is_read = false + AND deleted_at IS NULL`, + [randomUUID()], + ); + + return rows.map((plan) => plan["QUERY PLAN"]).join("\n"); + }); + + expect(plan).toContain("notifications_message_id_idx"); + }); + }); +}); From 3b841556dfba5ceeaa8aea7fa59fb5c02b1356d9 Mon Sep 17 00:00:00 2001 From: Luke Policinski Date: Mon, 28 Sep 2026 16:37:56 -0400 Subject: [PATCH 2/5] =?UTF-8?q?bug:=20chat=20self=20edit=20review=20fixes?= =?UTF-8?q?=20=E2=80=94=20edited=20preview=20on=20read=20rows,=20atomic=20?= =?UTF-8?q?draft=20move?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - updateChatMessagePreview rewrites every row for the message except a retracted (blank) one: read rows stay in the bell for a week and a recipient can restore a collapsed one, so both kept the text the author took back. `message <> ''` still stops an edit racing a delete from writing text back, including when it waits on the retraction's row lock - migrateLobbyMessages moves a draft's messages in one Lua script, so an edit or a self-delete in the draft room can no longer land on a copy that is thrown away (edit lost) or written back (deleted message returns) - lobby:edit answers an empty or non-string edit with chat:error invalid instead of leaving the client waiting - deleting the only message of a DM conversation takes it off both rails, so the recipient is not left with an empty tab from the sender - tests: the edit/delete race asserts the refusal code, the DM catch-up test deletes through the service, a retraction holding the row lock mid-edit --- src/chat/chat.gateway.spec.ts | 22 +++++++++- src/chat/chat.gateway.ts | 6 +-- src/chat/chat.service.spec.ts | 37 +++++++++++++--- src/chat/chat.service.ts | 49 ++++++++++++++-------- src/notifications/notifications.service.ts | 10 +++-- test/chat-direct-messages.spec.ts | 28 +++++++++++++ test/chat-redis-actions.spec.ts | 43 ++++++++++++++++--- 7 files changed, 158 insertions(+), 37 deletions(-) diff --git a/src/chat/chat.gateway.spec.ts b/src/chat/chat.gateway.spec.ts index bb5c3a87..cd65e765 100644 --- a/src/chat/chat.gateway.spec.ts +++ b/src/chat/chat.gateway.spec.ts @@ -507,8 +507,6 @@ describe("ChatGateway lobby:edit", () => { ["an unknown lobby type", { type: "global" }], ["a room id that is not a string", { id: 1 }], ["a missing message id", { messageId: undefined }], - ["a message that is not a string", { message: 5 }], - ["a message that is only whitespace", { message: " \n " }], ])("ignores %s", async (_, overrides) => { const socket = client(); @@ -518,6 +516,26 @@ describe("ChatGateway lobby:edit", () => { expect(socket.send).not.toHaveBeenCalled(); }); + it.each([ + ["only whitespace", " \n "], + ["not a string", 5], + ])("answers an edit that is %s with invalid", async (_, message) => { + const socket = client(); + + await gateway.editMessage( + edit({ message, requestId: "r-4" }) as any, + socket, + ); + + expect(chat.editMessage).not.toHaveBeenCalled(); + expect(sent(socket)).toEqual([ + { + event: "chat:error", + data: { code: ChatErrorCode.Invalid, action: "edit", requestId: "r-4" }, + }, + ]); + }); + it("ignores a missing payload", async () => { await gateway.editMessage(undefined as any, client()); diff --git a/src/chat/chat.gateway.ts b/src/chat/chat.gateway.ts index 59df2690..db281aab 100644 --- a/src/chat/chat.gateway.ts +++ b/src/chat/chat.gateway.ts @@ -224,10 +224,10 @@ export class ChatGateway { const parsed = ChatService.messageText(data.message); + // Unlike a send, clearing the box is an ordinary thing to do to an edit, + // so the client is told rather than left waiting. if ("error" in parsed) { - if (parsed.error === ChatErrorCode.TooLong) { - this.sendError(client, "edit", parsed.error, requestId); - } + this.sendError(client, "edit", parsed.error, requestId); return; } diff --git a/src/chat/chat.service.spec.ts b/src/chat/chat.service.spec.ts index 91339fae..d0aa4a20 100644 --- a/src/chat/chat.service.spec.ts +++ b/src/chat/chat.service.spec.ts @@ -1948,16 +1948,43 @@ describe("ChatService direct messages", () => { }); it("retracts them when a direct message was deleted meanwhile", async () => { - const messageId = await say( - ChatLobbyType.Direct, - directRoomId(ME, FRIEND), - ); + const room = directRoomId(ME, FRIEND); + notifications.notifyPlayers.mockImplementationOnce(async () => { + const insert = queries.find(({ sql }) => + sql.includes("INSERT INTO public.direct_messages"), + ); + directMessage = { + id: insert.bindings[0], + roomId: room, + author: ME, + message: "typo", + open: true, + editedAt: null, + }; + + await expect( + service.deleteMessage( + ChatLobbyType.Direct, + room, + insert.bindings[0], + me(), + ), + ).resolves.toEqual({ deleted: true }); + }); + + const messageId = await say(ChatLobbyType.Direct, room); await flush(); await flush(); - expect(notifications.retractChatMessage).toHaveBeenCalledWith( + expect(notifications.retractChatMessage).toHaveBeenCalledTimes(2); + expect(notifications.retractChatMessage).toHaveBeenLastCalledWith( messageId, ); + expect( + notifications.retractChatMessage.mock.invocationCallOrder[1], + ).toBeGreaterThan( + notifications.collapseOlderUnread.mock.invocationCallOrder[0], + ); }); it("gives them a direct message's edited text", async () => { diff --git a/src/chat/chat.service.ts b/src/chat/chat.service.ts index 6c683b11..847b4bf8 100644 --- a/src/chat/chat.service.ts +++ b/src/chat/chat.service.ts @@ -78,6 +78,19 @@ export class ChatService { return 1 `; + // 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. + 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]) + end + redis.call('DEL', KEYS[1]) + return #messages / 2 + `; + // 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. @@ -991,6 +1004,17 @@ export class ChatService { }; } + // Deleting the first message of a conversation would otherwise leave an + // empty tab on the recipient's rail, telling them something was sent. + await this.postgres.query( + `DELETE FROM public.direct_conversations + WHERE room_id = $1 + AND NOT EXISTS ( + SELECT 1 FROM public.direct_messages WHERE room_id = $1 + )`, + [roomId], + ); + void this.to(ChatLobbyType.Direct, roomId, "deleted", { id: messageId }); await this.retractNotifications(ChatLobbyType.Direct, roomId, messageId); @@ -2302,28 +2326,19 @@ export class ChatService { toType: ChatLobbyType, toId: string, ) { - const fromKey = `chat_${fromType}_${fromId}`; const toKey = `chat_${toType}_${toId}`; - const messagesObject = await this.redis.hgetall(fromKey); - - for (const [field, message] of Object.entries(messagesObject)) { - await this.redis.hset(toKey, field, message); - await this.redis.sendCommand( - new Redis.Command("HEXPIRE", [ - toKey, - this.ttlFor(toType), - "FIELDS", - 1, - field, - ]), - ); - } + const moved = await this.redis.eval( + ChatService.MOVE_ROOM_MESSAGES_SCRIPT, + 2, + `chat_${fromType}_${fromId}`, + toKey, + this.ttlFor(toType), + ); - await this.redis.del(fromKey); await this.removeLobby(fromType, fromId); - if (Object.keys(messagesObject).length === 0) { + if (moved === 0) { return; } diff --git a/src/notifications/notifications.service.ts b/src/notifications/notifications.service.ts index 14fcc1be..d9b62486 100644 --- a/src/notifications/notifications.service.ts +++ b/src/notifications/notifications.service.ts @@ -789,16 +789,18 @@ export class NotificationsService { ); } - // `deleted_at IS NULL` is what keeps an edit racing a delete from writing - // text back into a row retractChatMessage has already blanked. + // Read and collapsed rows too: the bell keeps read rows on show, and a + // recipient can restore a collapsed one, so either would otherwise keep the + // text the author took back. A preview is never empty, which is what tells a + // row retractChatMessage blanked apart -- an edit racing a delete must not + // write text back into it. async updateChatMessagePreview(messageId: string, preview: string) { await this.postgres.query( `UPDATE public.notifications SET message = $2 WHERE data->>'messageId' = $1 AND type IN ('ChatMessage', 'MatchChatMessage') - AND is_read = false - AND deleted_at IS NULL`, + AND message <> ''`, [messageId, preview], ); } diff --git a/test/chat-direct-messages.spec.ts b/test/chat-direct-messages.spec.ts index ad53c8df..a4fbd5a2 100644 --- a/test/chat-direct-messages.spec.ts +++ b/test/chat-direct-messages.spec.ts @@ -603,6 +603,34 @@ describe("direct messages (SQL-driven)", () => { expect(count).toBe("0"); }); + it("takes a conversation off both rails once its only message is deleted", async () => { + const me = await fx.player(); + const friend = await fx.player(); + const room = directRoomId(me, friend); + const id = await sent(room, me); + + await chat.deleteMessage(ChatLobbyType.Direct, room, id, as(me)); + + expect(await chat.getDirectConversations(as(friend))).toEqual([]); + expect(await chat.getDirectConversations(as(me))).toEqual([]); + }); + + it("keeps a conversation that still has messages", async () => { + const me = await fx.player(); + const friend = await fx.player(); + const room = directRoomId(me, friend); + await sent(room, friend, "hello"); + const id = await sent(room, me); + + await chat.deleteMessage(ChatLobbyType.Direct, room, id, as(me)); + + expect( + (await chat.getDirectConversations(as(friend))).map( + ({ roomId }) => roomId, + ), + ).toEqual([room]); + }); + it("answers not_found once it is gone", async () => { const me = await fx.player(); const friend = await fx.player(); diff --git a/test/chat-redis-actions.spec.ts b/test/chat-redis-actions.spec.ts index 1f121a04..6be4799e 100644 --- a/test/chat-redis-actions.spec.ts +++ b/test/chat-redis-actions.spec.ts @@ -280,7 +280,9 @@ describe("chat edits and self deletes (SQL-driven)", () => { ]); expect(deleted).toEqual({ deleted: true }); - expect([true, false]).toContain(edited.edited); + if (edited.edited === false) { + expect(edited.code).toBe(ChatErrorCode.NotFound); + } expect(await redis.hexists(key(matchId), id)).toBe(0); } @@ -345,6 +347,12 @@ describe("chat edits and self deletes (SQL-driven)", () => { edited_at: result.edited ? result.edited_at : undefined, }), ]); + expect(await expiresAt(matchId, id)).toBeGreaterThan(Date.now()); + expect(await redis.exists(`chat_draft_${draftId}`)).toBe(0); + + await expect( + chat.editMessage(ChatLobbyType.Draft, draftId, id, user, "again"), + ).resolves.toEqual({ edited: false, code: ChatErrorCode.NotFound }); }); it("audits an author deleting their own message", async () => { @@ -420,7 +428,9 @@ describe("chat edits and self deletes (SQL-driven)", () => { expect(await text(unread)).toBe("<b>fixed</b>"); }); - it("rewrites only unread, live rows for that message", async () => { + // A read row stays in the bell, and a collapsed one can be restored by + // its recipient, so neither may keep the text the author took back. + it("rewrites every row for that message, read or collapsed", async () => { const reader = await fx.player("Reader"); const messageId = randomUUID(); const unread = await row(reader, messageId); @@ -431,11 +441,33 @@ describe("chat edits and self deletes (SQL-driven)", () => { await notifications.updateChatMessagePreview(messageId, "fixed"); expect(await text(unread)).toBe("fixed"); - expect(await text(read)).toBe("typo"); - expect(await text(collapsed)).toBe("typo"); + expect(await text(read)).toBe("fixed"); + expect(await text(collapsed)).toBe("fixed"); expect(await text(other)).toBe("typo"); }); + it("leaves a row retracted mid-edit blank", async () => { + const reader = await fx.player("Reader"); + const messageId = randomUUID(); + const id = await row(reader, messageId); + let pending: Promise = Promise.resolve(); + + await postgres.transaction(async (client) => { + await client.query( + `UPDATE notifications SET deleted_at = now(), message = '' + WHERE id = $1::uuid`, + [id], + ); + + pending = notifications.updateChatMessagePreview(messageId, "fixed"); + await new Promise((resolve) => setTimeout(resolve, 200)); + }); + + await pending; + + expect(await text(id)).toBe(""); + }); + it("leaves a retracted row blank when the edit lands after the delete", async () => { const reader = await fx.player("Reader"); const messageId = randomUUID(); @@ -455,8 +487,7 @@ describe("chat edits and self deletes (SQL-driven)", () => { `EXPLAIN UPDATE notifications SET message = 'x' WHERE data->>'messageId' = $1 AND type IN ('ChatMessage', 'MatchChatMessage') - AND is_read = false - AND deleted_at IS NULL`, + AND message <> ''`, [randomUUID()], ); From 6fef7171cf0f6f58d8f0af5208951ab93e64d6ae Mon Sep 17 00:00:00 2001 From: Luke Policinski Date: Mon, 28 Sep 2026 16:46:56 -0400 Subject: [PATCH 3/5] feature: audit chat edits the way deletions are audited An author could post abuse and edit it into something harmless within the 10 minute window, and nothing of the original survived. Editing and then deleting kept only the edited text in chat_message_deletions. - chat_message_edits (message, room, author, previous and new text, when the message was sent, when it was edited), one row per group-room edit. DMs are never audited, the same as deletions. - The row is written before the redis compare-and-set, the way a deletion is audited before its HDEL, so a failed audit write refuses the edit. A swap that finds the message gone or changed deletes its row again; a swap that fails outright keeps it, since it may have applied. - The swap leaves a short-lived receipt keyed on its audit row, so an ioredis resend after a lost reply answers 1 even if the message was deleted or edited again in between, instead of reading as a failed swap and discarding the audit of an edit that happened. - Hasura: moderator reads every row outside the organizers' room, match_organizer and up read everything; no insert/update/delete. --- .../tables/public_chat_message_edits.yaml | 40 +++ .../databases/default/tables/tables.yaml | 1 + .../1888000000250_chat_message_edits/down.sql | 1 + .../1888000000250_chat_message_edits/up.sql | 20 ++ src/chat/chat.service.spec.ts | 177 ++++++++++++ src/chat/chat.service.ts | 89 ++++++- test/chat-direct-messages.spec.ts | 5 + test/chat-redis-actions.spec.ts | 252 +++++++++++++++++- 8 files changed, 579 insertions(+), 6 deletions(-) create mode 100644 hasura/metadata/databases/default/tables/public_chat_message_edits.yaml create mode 100644 hasura/migrations/default/1888000000250_chat_message_edits/down.sql create mode 100644 hasura/migrations/default/1888000000250_chat_message_edits/up.sql diff --git a/hasura/metadata/databases/default/tables/public_chat_message_edits.yaml b/hasura/metadata/databases/default/tables/public_chat_message_edits.yaml new file mode 100644 index 00000000..5cf0c3c8 --- /dev/null +++ b/hasura/metadata/databases/default/tables/public_chat_message_edits.yaml @@ -0,0 +1,40 @@ +table: + name: chat_message_edits + schema: public +object_relationships: + - name: author + using: + foreign_key_constraint_on: author_steam_id +select_permissions: + - role: match_organizer + permission: + columns: + - id + - message_id + - room_type + - room_id + - author_steam_id + - previous_message + - new_message + - message_created_at + - edited_at + filter: {} + allow_aggregations: true + comment: What website chat said before its author edited it. Written only by the API. + - role: moderator + permission: + columns: + - id + - message_id + - room_type + - room_id + - author_steam_id + - previous_message + - new_message + - message_created_at + - edited_at + filter: + room_type: + _neq: organizers + allow_aggregations: true + comment: The organizers' room is closed to moderators, and so is its evidence. diff --git a/hasura/metadata/databases/default/tables/tables.yaml b/hasura/metadata/databases/default/tables/tables.yaml index 32afb034..07d8efd5 100644 --- a/hasura/metadata/databases/default/tables/tables.yaml +++ b/hasura/metadata/databases/default/tables/tables.yaml @@ -6,6 +6,7 @@ - "!include public_awards.yaml" - "!include public_broadcast_huds.yaml" - "!include public_chat_message_deletions.yaml" +- "!include public_chat_message_edits.yaml" - "!include public_chat_read_state.yaml" - "!include public_clip_render_jobs.yaml" - "!include public_custom_pages.yaml" diff --git a/hasura/migrations/default/1888000000250_chat_message_edits/down.sql b/hasura/migrations/default/1888000000250_chat_message_edits/down.sql new file mode 100644 index 00000000..cededa6c --- /dev/null +++ b/hasura/migrations/default/1888000000250_chat_message_edits/down.sql @@ -0,0 +1 @@ +DROP TABLE IF EXISTS public.chat_message_edits; diff --git a/hasura/migrations/default/1888000000250_chat_message_edits/up.sql b/hasura/migrations/default/1888000000250_chat_message_edits/up.sql new file mode 100644 index 00000000..337c38a3 --- /dev/null +++ b/hasura/migrations/default/1888000000250_chat_message_edits/up.sql @@ -0,0 +1,20 @@ +CREATE TABLE IF NOT EXISTS public.chat_message_edits ( + id uuid NOT NULL DEFAULT gen_random_uuid(), + message_id uuid NOT NULL, + room_type text NOT NULL, + room_id text NOT NULL, + author_steam_id bigint REFERENCES public.players (steam_id) + ON UPDATE CASCADE ON DELETE SET NULL, + previous_message text NOT NULL, + new_message text NOT NULL, + message_created_at timestamptz, + edited_at timestamptz NOT NULL DEFAULT now(), + + PRIMARY KEY (id) +); + +CREATE INDEX IF NOT EXISTS chat_message_edits_message_id_idx + ON public.chat_message_edits (message_id); + +CREATE INDEX IF NOT EXISTS chat_message_edits_author_idx + ON public.chat_message_edits (author_steam_id, edited_at DESC); diff --git a/src/chat/chat.service.spec.ts b/src/chat/chat.service.spec.ts index d0aa4a20..9c5b1bd4 100644 --- a/src/chat/chat.service.spec.ts +++ b/src/chat/chat.service.spec.ts @@ -38,6 +38,7 @@ describe("ChatService direct messages", () => { let queries: Array<{ sql: string; bindings: any[] }>; let gagged: boolean; let audited: boolean; + let editAuditIds: string[]; // The one direct message the fake database holds, if a test put one there. let directMessage: | { @@ -69,6 +70,12 @@ describe("ChatService direct messages", () => { return [{ deleted: audited }]; } + if (sql.includes("INSERT INTO public.chat_message_edits")) { + const id = `edit-audit-${editAuditIds.length + 1}`; + editAuditIds.push(id); + return [{ id }]; + } + if (sql.includes("AS open")) { return directMessage?.id === bindings[0] && directMessage.roomId === bindings[1] @@ -340,6 +347,7 @@ describe("ChatService direct messages", () => { queries = []; gagged = false; audited = false; + editAuditIds = []; rcon.send.mockResolvedValue(undefined); rcon.connect.mockResolvedValue(rcon); @@ -1418,6 +1426,18 @@ describe("ChatService direct messages", () => { sql.includes("INSERT INTO public.chat_message_deletions"), ); + const editAudits = () => + queries.filter(({ sql }) => + sql.includes("INSERT INTO public.chat_message_edits"), + ); + + const discardedEditAudits = () => + queries + .filter(({ sql }) => + sql.includes("DELETE FROM public.chat_message_edits"), + ) + .map(({ bindings }) => bindings[0]); + const flush = () => new Promise((resolve) => setImmediate(resolve)); beforeEach(() => { @@ -1439,6 +1459,7 @@ describe("ChatService direct messages", () => { script: string, _keys: number, key: string, + _receipt: string, field: string, expected: string, next: string, @@ -1528,6 +1549,145 @@ describe("ChatService direct messages", () => { expect(rcon.connect).not.toHaveBeenCalled(); }); + it("keeps what the message said before the edit", async () => { + store(); + const { timestamp } = current(); + + const result = await edit(); + + expect(editAudits().map(({ bindings }) => bindings)).toEqual([ + [ + MESSAGE_ID, + "match", + "m-1", + ME, + "typo", + "fixed", + timestamp, + result.edited ? result.edited_at : undefined, + ], + ]); + expect(discardedEditAudits()).toEqual([]); + }); + + it("writes the audit row before the message changes", async () => { + store(); + + await edit(); + + const auditCall = postgres.query.mock.calls.findIndex(([sql]) => + sql.includes("INSERT INTO public.chat_message_edits"), + ); + const swapCall = redis.eval.mock.calls.findIndex(([script]) => + script.includes("HPEXPIRETIME"), + ); + + expect(auditCall).toBeGreaterThanOrEqual(0); + expect(postgres.query.mock.invocationCallOrder[auditCall]).toBeLessThan( + redis.eval.mock.invocationCallOrder[swapCall], + ); + }); + + it("leaves the message as it was when the audit row cannot be written", async () => { + store(); + const database = postgres.query.getMockImplementation(); + postgres.query.mockImplementation(async (sql, bindings) => { + if (sql.includes("INSERT INTO public.chat_message_edits")) { + throw new Error("database down"); + } + return database(sql, bindings); + }); + + try { + await expect(edit()).rejects.toThrow("database down"); + } finally { + postgres.query.mockImplementation(database); + } + + await flush(); + + expect(current().message).toBe("typo"); + expect( + redis.eval.mock.calls.some(([script]) => + script.includes("HPEXPIRETIME"), + ), + ).toBe(false); + expect(broadcasts("lobby:match:m-1:edited")).toEqual([]); + }); + + it("keeps the audit row when the swap fails outright, since it may have applied", async () => { + store(); + redis.eval.mockRejectedValueOnce(new Error("connection reset")); + + await expect(edit()).rejects.toThrow("connection reset"); + + expect(editAudits()).toHaveLength(1); + expect(discardedEditAudits()).toEqual([]); + }); + + it("hands the swap a receipt named for its own audit row", async () => { + store(); + + await edit(); + + const swap = redis.eval.mock.calls.find(([script]) => + script.includes("HPEXPIRETIME"), + ); + + expect(swap.slice(1, 4)).toEqual([ + 2, + "chat_match_m-1", + "chat_edit_applied:edit-audit-1", + ]); + }); + + it("still retries when an audit row it no longer needs cannot be discarded", async () => { + store(); + redis.hget.mockImplementationOnce( + async (key: string, field: string) => { + const raw = stored[key][field]; + stored[key][field] = JSON.stringify({ + ...JSON.parse(raw), + message: "other tab", + }); + return raw; + }, + ); + const database = postgres.query.getMockImplementation(); + postgres.query.mockImplementation(async (sql, bindings) => { + if (sql.includes("DELETE FROM public.chat_message_edits")) { + throw new Error("database blip"); + } + return database(sql, bindings); + }); + + try { + await expect(edit()).resolves.toMatchObject({ edited: true }); + } finally { + postgres.query.mockImplementation(database); + } + + expect(current().message).toBe("fixed"); + expect(logger.warn).toHaveBeenCalledWith( + expect.stringContaining("unable to discard"), + expect.any(Error), + ); + }); + + it("audits nothing for an edit it refuses", async () => { + store({ timestamp: ago(WINDOW + 6_000) }); + await edit(); + + store(); + gagged = true; + await edit(); + + store({ source: "game" }); + await edit(); + + expect(editAudits()).toEqual([]); + }); + it.each([ ["another player", { steam_id: FRIEND }], ["an administrator", { steam_id: FRIEND, role: "administrator" }], @@ -1636,6 +1796,8 @@ describe("ChatService direct messages", () => { }); expect(current()).toBeNull(); expect(broadcasts("lobby:match:m-1:edited")).toEqual([]); + expect(discardedEditAudits()).toEqual(editAuditIds); + expect(editAuditIds).toHaveLength(1); }); it("applies an edit on top of one that landed between the read and the write", async () => { @@ -1654,6 +1816,11 @@ describe("ChatService direct messages", () => { await expect(edit()).resolves.toMatchObject({ edited: true }); expect(current().message).toBe("fixed"); + expect(editAudits().map(({ bindings }) => bindings[4])).toEqual([ + "typo", + "other tab", + ]); + expect(discardedEditAudits()).toEqual(["edit-audit-1"]); }); it("lets the author delete their own message, and audits it", async () => { @@ -1788,6 +1955,16 @@ describe("ChatService direct messages", () => { ); }); + it("never audits the edit of a direct message", async () => { + hold(); + + await expect(editDirect()).resolves.toMatchObject({ edited: true }); + + expect( + queries.some(({ sql }) => sql.includes("chat_message_edits")), + ).toBe(false); + }); + it("is not held back by a gag", async () => { hold(); gagged = true; diff --git a/src/chat/chat.service.ts b/src/chat/chat.service.ts index 847b4bf8..840c0d9c 100644 --- a/src/chat/chat.service.ts +++ b/src/chat/chat.service.ts @@ -66,7 +66,16 @@ export class ChatService { // back: an edit never extends a message's life. Comparing against the value // that was read and checked keeps an edit from bringing back a message that // expired or was deleted in the meantime, or from overwriting another edit. + // + // KEYS[2] is a receipt for this one attempt. ioredis resends a command whose + // reply was lost to a reconnect, seconds later, by which time the message may + // have been deleted or edited again; the receipt is what tells that resend it + // already applied, rather than reading as a failed swap and discarding the + // audit row of an edit that happened. private static readonly EDIT_ROOM_MESSAGE_SCRIPT = ` + if redis.call('EXISTS', KEYS[2]) == 1 then + return 1 + end if redis.call('HGET', KEYS[1], ARGV[1]) ~= ARGV[2] then return 0 end @@ -75,9 +84,12 @@ export class ChatService { if expiresAt > 0 then redis.call('HPEXPIREAT', KEYS[1], expiresAt, 'FIELDS', 1, ARGV[1]) end + redis.call('SET', KEYS[2], '1', 'PX', ARGV[4]) return 1 `; + 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. @@ -908,18 +920,34 @@ export class ChatService { const editedAt = new Date().toISOString(); + // Written before the swap, the way a deletion is audited before its + // HDEL, so no failure part way through can replace what was said + // without keeping it. A swap that does not apply takes its row back out. + const auditId = await this.recordEdit( + type, + id, + messageId, + message, + text, + editedAt, + ); + const swapped = await this.redis.eval( ChatService.EDIT_ROOM_MESSAGE_SCRIPT, - 1, + 2, messageKey, + `chat_edit_applied:${auditId}`, messageId, raw, JSON.stringify({ ...message, message: text, edited_at: editedAt }), + ChatService.EDIT_RECEIPT_TTL_MS, ); if (swapped === 1) { return await this.announceEdit(type, id, messageId, text, editedAt); } + + await this.discardEdit(auditId); } this.logger.warn( @@ -1116,8 +1144,6 @@ export class ChatService { message: ChatMessage, deletedBy: User, ) { - const createdAt = new Date(message.timestamp); - await this.postgres.query( `INSERT INTO public.chat_message_deletions (message_id, room_type, room_id, author_steam_id, message, @@ -1133,13 +1159,68 @@ export class ChatService { id, ChatService.authorSteamId(message), String(message.message ?? ""), - Number.isNaN(createdAt.getTime()) ? null : createdAt.toISOString(), + ChatService.messageCreatedAt(message), message.source ?? null, deletedBy.steam_id, ], ); } + // edited_at is the one the message itself now carries, so a row can be + // matched to the edit clients were shown. + private async recordEdit( + type: ChatLobbyType, + id: string, + messageId: string, + message: ChatMessage, + text: string, + editedAt: string, + ): Promise { + const [row] = await this.postgres.query>( + `INSERT INTO public.chat_message_edits + (message_id, room_type, room_id, author_steam_id, + previous_message, new_message, message_created_at, edited_at) + SELECT $1::uuid, $2, $3, + (SELECT steam_id FROM public.players + WHERE steam_id = $4::bigint), + $5, $6, $7::timestamptz, $8::timestamptz + RETURNING id::text AS id`, + [ + messageId, + type, + id, + ChatService.authorSteamId(message), + String(message.message ?? ""), + text, + ChatService.messageCreatedAt(message), + editedAt, + ], + ); + + return row.id; + } + + // Only for a swap that found the message gone or changed, so the edit never + // happened. A swap that failed outright keeps its row: it may have applied. + private async discardEdit(auditId: string) { + await this.postgres + .query(`DELETE FROM public.chat_message_edits WHERE id = $1::uuid`, [ + auditId, + ]) + .catch((error) => { + this.logger.warn( + `unable to discard the audit row of an edit that did not apply`, + error, + ); + }); + } + + private static messageCreatedAt(message: ChatMessage): string | null { + const sentAt = new Date(message.timestamp); + + return Number.isNaN(sentAt.getTime()) ? null : sentAt.toISOString(); + } + // A steam id stored as a JSON number has already been rounded by JSON.parse // onto some other account, and would pin the message on the wrong player. private static authorSteamId(message: ChatMessage): string | null { diff --git a/test/chat-direct-messages.spec.ts b/test/chat-direct-messages.spec.ts index a4fbd5a2..39b615ab 100644 --- a/test/chat-direct-messages.spec.ts +++ b/test/chat-direct-messages.spec.ts @@ -488,6 +488,11 @@ describe("direct messages (SQL-driven)", () => { message: "fixed", edited_at: after.edited_at.toISOString(), }); + + const [{ count }] = await postgres.query>( + `SELECT count(*)::text AS count FROM chat_message_edits`, + ); + expect(count).toBe("0"); }); it("shows the edit in history after a reload", async () => { diff --git a/test/chat-redis-actions.spec.ts b/test/chat-redis-actions.spec.ts index 6be4799e..1b99cd31 100644 --- a/test/chat-redis-actions.spec.ts +++ b/test/chat-redis-actions.spec.ts @@ -1,4 +1,6 @@ import { randomUUID } from "crypto"; +import { readFileSync } from "fs"; +import { join } from "path"; import IORedis, { Redis } from "ioredis"; import { GenericContainer, StartedTestContainer } from "testcontainers"; import { PostgresService } from "./../src/postgres/postgres.service"; @@ -100,6 +102,7 @@ describe("chat edits and self deletes (SQL-driven)", () => { jest.clearAllMocks(); await redis.flushall(); await postgres.query("DELETE FROM chat_message_deletions"); + await postgres.query("DELETE FROM chat_message_edits"); await postgres.query("DELETE FROM player_sanctions"); await postgres.query("DELETE FROM notifications"); await postgres.query("DELETE FROM players"); @@ -253,23 +256,45 @@ describe("chat edits and self deletes (SQL-driven)", () => { await redis.call("HPEXPIREAT", key(matchId), at, "FIELDS", 1, id); const script = (ChatService as any).EDIT_ROOM_MESSAGE_SCRIPT; + const receipt = `chat_edit_applied:${randomUUID()}`; await expect( - redis.eval(script, 1, key(matchId), id, "stale", "replacement"), + redis.eval( + script, + 2, + key(matchId), + receipt, + id, + "stale", + "replacement", + 60_000, + ), ).resolves.toBe(0); await expect( - redis.eval(script, 1, key(matchId), randomUUID(), "", "replacement"), + redis.eval( + script, + 2, + key(matchId), + receipt, + randomUUID(), + "", + "replacement", + 60_000, + ), ).resolves.toBe(0); expect(await redis.hget(key(matchId), id)).toBe(current); expect(await expiresAt(matchId, id)).toBe(at); expect(await redis.hlen(key(matchId))).toBe(1); + expect(await redis.exists(receipt)).toBe(0); }); }); it("never resurrects a message when an edit races its deletion", async () => { const user = await author(); + const applied: string[] = []; + for (let round = 0; round < 25; round++) { const matchId = randomUUID(); const id = await place(matchId, user); @@ -282,6 +307,8 @@ describe("chat edits and self deletes (SQL-driven)", () => { expect(deleted).toEqual({ deleted: true }); if (edited.edited === false) { expect(edited.code).toBe(ChatErrorCode.NotFound); + } else { + applied.push(id); } expect(await redis.hexists(key(matchId), id)).toBe(0); } @@ -293,6 +320,14 @@ describe("chat edits and self deletes (SQL-driven)", () => { ); expect(count).toBe("25"); + + const audited = await postgres.query>( + `SELECT message_id::text AS message_id FROM chat_message_edits`, + ); + + expect(audited.map(({ message_id }) => message_id).sort()).toEqual( + applied.sort(), + ); }); it("stores the edit exactly as JSON.stringify writes it", async () => { @@ -378,6 +413,219 @@ describe("chat edits and self deletes (SQL-driven)", () => { ]); }); + describe("the edit audit", () => { + const up = readFileSync( + join( + __dirname, + "../hasura/migrations/default/1888000000250_chat_message_edits/up.sql", + ), + "utf8", + ); + + const down = readFileSync( + join( + __dirname, + "../hasura/migrations/default/1888000000250_chat_message_edits/down.sql", + ), + "utf8", + ); + + const edits = (messageId: string) => + postgres.query< + Array<{ + room_type: string; + room_id: string; + author: string | null; + previous_message: string; + new_message: string; + message_created_at: Date | null; + edited_at: Date; + }> + >( + `SELECT room_type, room_id, author_steam_id::text AS author, + previous_message, new_message, message_created_at, edited_at + FROM chat_message_edits + WHERE message_id = $1::uuid + ORDER BY edited_at`, + [messageId], + ); + + it("keeps what the message said before each edit", async () => { + const user = await author(); + const matchId = randomUUID(); + const id = await place(matchId, user); + const { timestamp } = JSON.parse(await redis.hget(key(matchId), id)); + + const first = await edit(matchId, id, user, "fixed"); + await new Promise((resolve) => setTimeout(resolve, 5)); + const second = await edit(matchId, id, user, "fixed again"); + + expect(await edits(id)).toEqual([ + { + room_type: "match", + room_id: matchId, + author: user.steam_id, + previous_message: "typo", + new_message: "fixed", + message_created_at: new Date(timestamp), + edited_at: new Date(first.edited ? first.edited_at : 0), + }, + { + room_type: "match", + room_id: matchId, + author: user.steam_id, + previous_message: "fixed", + new_message: "fixed again", + message_created_at: new Date(timestamp), + edited_at: new Date(second.edited ? second.edited_at : 0), + }, + ]); + }); + + it("keeps the original when a message is edited and then deleted", async () => { + const user = await author(); + const matchId = randomUUID(); + const id = await place(matchId, user); + + await edit(matchId, id, user, "harmless"); + await chat.deleteMessage(ChatLobbyType.Match, matchId, id, user); + + const [deletion] = await postgres.query>( + `SELECT message FROM chat_message_deletions WHERE message_id = $1::uuid`, + [id], + ); + + expect(deletion.message).toBe("harmless"); + expect((await edits(id)).map((row) => row.previous_message)).toEqual([ + "typo", + ]); + }); + + it("records nothing for an edit that never applied", async () => { + const user = await author(); + const matchId = randomUUID(); + const id = await place(matchId, user); + + pauseFirstRead(() => redis.hdel(key(matchId), id)); + + await expect(edit(matchId, id, user)).resolves.toEqual({ + edited: false, + code: ChatErrorCode.NotFound, + }); + + expect(await edits(id)).toEqual([]); + }); + + it("keeps only the edit that applied when another lands in between", async () => { + const user = await author(); + const matchId = randomUUID(); + const id = await place(matchId, user); + + pauseFirstRead(() => + chat.editMessage(ChatLobbyType.Match, matchId, id, user, "other tab"), + ); + + await expect(edit(matchId, id, user, "mine")).resolves.toMatchObject({ + edited: true, + }); + + expect( + (await edits(id)).map(({ previous_message, new_message }) => [ + previous_message, + new_message, + ]), + ).toEqual([ + ["typo", "other tab"], + ["other tab", "mine"], + ]); + }); + + // The resend comes after ioredis reconnects, seconds later, so the message + // may already be gone or edited again by then. + it.each([ + [ + "deleted", + (matchId: string, id: string) => redis.hdel(key(matchId), id), + ], + [ + "edited again", + (matchId: string, id: string) => + redis.hset( + key(matchId), + id, + JSON.stringify({ id, message: "other" }), + ), + ], + ])( + "counts a resent swap that already applied as applied, though the message was since %s", + async (_, meanwhile) => { + const user = await author(); + const matchId = randomUUID(); + const id = await place(matchId, user); + const raw = await redis.hget(key(matchId), id); + const next = JSON.stringify({ ...JSON.parse(raw), message: "fixed" }); + const script = (ChatService as any).EDIT_ROOM_MESSAGE_SCRIPT; + const receipt = `chat_edit_applied:${randomUUID()}`; + const swap = () => + redis.eval(script, 2, key(matchId), receipt, id, raw, next, 60_000); + + await expect(swap()).resolves.toBe(1); + expect(await redis.hget(key(matchId), id)).toBe(next); + expect(await redis.pttl(receipt)).toBeGreaterThan(0); + + await meanwhile(matchId, id); + const since = await redis.hget(key(matchId), id); + + await expect(swap()).resolves.toBe(1); + expect(await redis.hget(key(matchId), id)).toBe(since); + }, + ); + + it("outlives the author's player row", async () => { + const user = await author(); + const matchId = randomUUID(); + const id = await place(matchId, user); + + await edit(matchId, id, user); + await postgres.query(`DELETE FROM players WHERE steam_id = $1::bigint`, [ + user.steam_id, + ]); + + expect(await edits(id)).toEqual([ + expect.objectContaining({ author: null, previous_message: "typo" }), + ]); + }); + + it("re-applies the migration cleanly and rolls back", async () => { + const table = async () => + ( + await postgres.query>( + `SELECT to_regclass('public.chat_message_edits')::text AS table`, + ) + )[0].table; + + await expect(postgres.query(up)).resolves.toBeDefined(); + await expect(postgres.query(up)).resolves.toBeDefined(); + await expect(postgres.query(down)).resolves.toBeDefined(); + expect(await table()).toBeNull(); + await expect(postgres.query(down)).resolves.toBeDefined(); + await expect(postgres.query(up)).resolves.toBeDefined(); + expect(await table()).toBe("chat_message_edits"); + + const indexes = await postgres.query>( + `SELECT indexname FROM pg_indexes + WHERE schemaname = 'public' AND tablename = 'chat_message_edits' + ORDER BY indexname`, + ); + + expect(indexes.map(({ indexname }) => indexname)).toEqual([ + "chat_message_edits_author_idx", + "chat_message_edits_message_id_idx", + "chat_message_edits_pkey", + ]); + }); + }); + describe("the bell's preview", () => { const row = async ( steamId: string, From bc1908f717341e0bff476a0f4cfec1fa5a0a45d2 Mon Sep 17 00:00:00 2001 From: Luke Policinski Date: Mon, 28 Sep 2026 18:19:49 -0400 Subject: [PATCH 4/5] bug: carry the server's edited_at on the edit ack --- src/chat/chat.gateway.spec.ts | 15 +++++++++++++-- src/chat/chat.gateway.ts | 8 ++++++-- 2 files changed, 19 insertions(+), 4 deletions(-) diff --git a/src/chat/chat.gateway.spec.ts b/src/chat/chat.gateway.spec.ts index cd65e765..4121b04b 100644 --- a/src/chat/chat.gateway.spec.ts +++ b/src/chat/chat.gateway.spec.ts @@ -579,7 +579,12 @@ describe("ChatGateway lobby:edit", () => { ]); }); - it("acks an edit under the requestId it came with", async () => { + it("acks an edit under the requestId it came with, with what the server stored", async () => { + chat.editMessage.mockResolvedValue({ + edited: true, + message: "fixed as stored", + edited_at: "2026-02-03T04:05:06.789Z", + }); const socket = client(); await gateway.editMessage(edit({ requestId: "r-2" }) as any, socket); @@ -587,7 +592,13 @@ describe("ChatGateway lobby:edit", () => { expect(sent(socket)).toEqual([ { event: "chat:ack", - data: { requestId: "r-2", messageId: MESSAGE_ID, action: "edit" }, + data: { + requestId: "r-2", + messageId: MESSAGE_ID, + action: "edit", + message: "fixed as stored", + edited_at: "2026-02-03T04:05:06.789Z", + }, }, ]); }); diff --git a/src/chat/chat.gateway.ts b/src/chat/chat.gateway.ts index db281aab..1a162e04 100644 --- a/src/chat/chat.gateway.ts +++ b/src/chat/chat.gateway.ts @@ -245,7 +245,10 @@ export class ChatGateway { } if (requestId) { - this.sendAck(client, "edit", requestId, data.messageId); + this.sendAck(client, "edit", requestId, data.messageId, { + message: result.message, + edited_at: result.edited_at, + }); } } @@ -279,11 +282,12 @@ export class ChatGateway { action: ChatAction, requestId: string, messageId: string, + extra: Record = {}, ) { client.send( JSON.stringify({ event: "chat:ack", - data: { requestId, messageId, action }, + data: { ...extra, requestId, messageId, action }, }), ); } From 5715f46f163dea884551606994985c58dddd44bd Mon Sep 17 00:00:00 2001 From: Luke Policinski Date: Mon, 28 Sep 2026 20:31:37 -0400 Subject: [PATCH 5/5] test: prove a draft room's move into its match is atomic With the move made of separate calls, an author's delete in the draft room landed on a copy that was written back into the match, and an edit landed on one that was thrown away. --- test/chat-redis-actions.spec.ts | 61 +++++++++++++++++++++++++++++++++ 1 file changed, 61 insertions(+) diff --git a/test/chat-redis-actions.spec.ts b/test/chat-redis-actions.spec.ts index 1b99cd31..8ca15907 100644 --- a/test/chat-redis-actions.spec.ts +++ b/test/chat-redis-actions.spec.ts @@ -390,6 +390,67 @@ describe("chat edits and self deletes (SQL-driven)", () => { ).resolves.toEqual({ edited: false, code: ChatErrorCode.NotFound }); }); + describe("an action in the draft room while it moves into the match", () => { + // Runs the action the first time the move writes into the match room from + // outside redis, which is the gap a move made of separate calls leaves. + const midMove = (action: () => Promise) => { + let running: Promise | undefined; + const write = redis.hset.bind(redis) as ( + ...args: any[] + ) => Promise; + + jest.spyOn(redis, "hset").mockImplementation((async (...args: any[]) => { + running ??= action(); + await running; + return write(...args); + }) as any); + + return () => (running ??= action()); + }; + + const move = (draftId: string, matchId: string) => + chat.migrateLobbyMessages( + ChatLobbyType.Draft, + draftId, + ChatLobbyType.Match, + matchId, + ); + + it("never brings back a message the author deleted", async () => { + const user = await author(); + const draftId = randomUUID(); + const matchId = randomUUID(); + const id = randomUUID(); + await redis.hset(`chat_draft_${draftId}`, id, written(id, user)); + + const deleting = midMove(() => + chat.deleteMessage(ChatLobbyType.Draft, draftId, id, user), + ); + await move(draftId, matchId); + const { deleted } = await deleting(); + + expect(await redis.hexists(key(matchId), id)).toBe(deleted ? 0 : 1); + }); + + it("never loses an edit the author was told had applied", async () => { + const user = await author(); + const draftId = randomUUID(); + const matchId = randomUUID(); + const id = randomUUID(); + await redis.hset(`chat_draft_${draftId}`, id, written(id, user)); + + const editing = midMove(() => + chat.editMessage(ChatLobbyType.Draft, draftId, id, user, "fixed"), + ); + await move(draftId, matchId); + const { edited } = await editing(); + + expect(JSON.parse(await redis.hget(key(matchId), id)).message).toBe( + edited ? "fixed" : "typo", + ); + }); + }); + it("audits an author deleting their own message", async () => { const user = await author(); const matchId = randomUUID();