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/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/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.gateway.spec.ts b/src/chat/chat.gateway.spec.ts index 730692d8..4121b04b 100644 --- a/src/chat/chat.gateway.spec.ts +++ b/src/chat/chat.gateway.spec.ts @@ -461,3 +461,179 @@ 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 }], + ])("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.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()); + + 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, 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); + + expect(sent(socket)).toEqual([ + { + event: "chat:ack", + data: { + requestId: "r-2", + messageId: MESSAGE_ID, + action: "edit", + message: "fixed as stored", + edited_at: "2026-02-03T04:05:06.789Z", + }, + }, + ]); + }); + + 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..1a162e04 100644 --- a/src/chat/chat.gateway.ts +++ b/src/chat/chat.gateway.ts @@ -195,6 +195,63 @@ 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); + + // 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) { + 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, { + message: result.message, + edited_at: result.edited_at, + }); + } + } + private static isLobbyType(value: unknown): value is ChatLobbyType { return Object.values(ChatLobbyType).includes(value as ChatLobbyType); } @@ -225,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 }, }), ); } diff --git a/src/chat/chat.service.spec.ts b/src/chat/chat.service.spec.ts index 196753cd..9c5b1bd4 100644 --- a/src/chat/chat.service.spec.ts +++ b/src/chat/chat.service.spec.ts @@ -38,6 +38,26 @@ 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: + | { + 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 +70,59 @@ 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] + ? [{ 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 +132,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 +329,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"]; @@ -270,6 +347,7 @@ describe("ChatService direct messages", () => { queries = []; gagged = false; audited = false; + editAuditIds = []; rcon.send.mockResolvedValue(undefined); rcon.connect.mockResolvedValue(rcon); @@ -1178,8 +1256,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 +1276,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 +1377,824 @@ 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 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(() => { + 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, + _receipt: 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("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" }], + ])("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([]); + 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 () => { + 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"); + 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 () => { + 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("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; + + 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 room = directRoomId(ME, FRIEND); + notifications.notifyPlayers.mockImplementationOnce(async () => { + const insert = queries.find(({ sql }) => + sql.includes("INSERT INTO public.direct_messages"), + ); + directMessage = { + id: insert.bindings[0], + roomId: room, + author: ME, + message: "typo", + open: true, + editedAt: null, + }; + + await expect( + service.deleteMessage( + ChatLobbyType.Direct, + room, + insert.bindings[0], + me(), + ), + ).resolves.toEqual({ deleted: true }); + }); + + const messageId = await say(ChatLobbyType.Direct, room); + await flush(); + await flush(); + + expect(notifications.retractChatMessage).toHaveBeenCalledTimes(2); + expect(notifications.retractChatMessage).toHaveBeenLastCalledWith( + messageId, + ); + expect( + notifications.retractChatMessage.mock.invocationCallOrder[1], + ).toBeGreaterThan( + notifications.collapseOlderUnread.mock.invocationCallOrder[0], + ); + }); + + it("gives them a direct message's edited text", async () => { + notifications.notifyPlayers.mockImplementationOnce(async () => { + const insert = queries.find(({ sql }) => + sql.includes("INSERT INTO public.direct_messages"), + ); + directMessage = { + id: insert.bindings[0], + roomId: insert.bindings[1], + author: ME, + message: "fixed", + open: true, + editedAt: new Date(), + }; + }); + + const messageId = await say( + ChatLobbyType.Direct, + directRoomId(ME, FRIEND), + ); + await flush(); + await flush(); + + expect(notifications.updateChatMessagePreview).toHaveBeenCalledWith( + messageId, + "fixed", + ); + expect(notifications.retractChatMessage).not.toHaveBeenCalled(); + }); + }); + }); + describe("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..840c0d9c 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,64 @@ 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. + // + // 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 + 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 + 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. + 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. + 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 +753,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 +768,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 +776,365 @@ 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(); + + // 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, + 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( + `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, + }; + } + + // 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); + + 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( @@ -757,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, @@ -774,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 { @@ -839,9 +1279,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,12 +1302,65 @@ 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 { @@ -1395,6 +1886,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 +1894,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 +1912,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 +2075,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); @@ -1905,28 +2407,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/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..d9b62486 100644 --- a/src/notifications/notifications.service.ts +++ b/src/notifications/notifications.service.ts @@ -789,6 +789,22 @@ export class NotificationsService { ); } + // Read and collapsed rows too: the bell keeps read rows on show, and a + // recipient can restore a collapsed one, so either would otherwise keep the + // text the author took back. A preview is never empty, which is what tells a + // row retractChatMessage blanked apart -- an edit racing a delete must not + // write text back into it. + async updateChatMessagePreview(messageId: string, preview: string) { + await this.postgres.query( + `UPDATE public.notifications + SET message = $2 + WHERE data->>'messageId' = $1 + AND type IN ('ChatMessage', 'MatchChatMessage') + AND message <> ''`, + [messageId, preview], + ); + } + // Opening a conversation should clear its badge everywhere, not just in the // tab that was open. async markConversationRead( diff --git a/test/chat-direct-messages.spec.ts b/test/chat-direct-messages.spec.ts index 94869e44..39b615ab 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,275 @@ 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(), + }); + + 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 () => { + 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("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(); + 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..8ca15907 --- /dev/null +++ b/test/chat-redis-actions.spec.ts @@ -0,0 +1,809 @@ +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"; +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 chat_message_edits"); + 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; + const receipt = `chat_edit_applied:${randomUUID()}`; + + await expect( + redis.eval( + script, + 2, + key(matchId), + receipt, + id, + "stale", + "replacement", + 60_000, + ), + ).resolves.toBe(0); + await expect( + 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); + + const [edited, deleted] = await Promise.all([ + edit(matchId, id, user), + chat.deleteMessage(ChatLobbyType.Match, matchId, id, user), + ]); + + 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); + } + + 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"); + + 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 () => { + 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, + }), + ]); + 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 }); + }); + + 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(); + 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 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, + messageId: string, + overrides: { + is_read?: boolean; + deleted?: boolean; + message?: string; + } = {}, + ) => { + const [inserted] = await postgres.query>( + `INSERT INTO notifications + (type, title, message, role, steam_id, entity_id, data, + is_read, deleted_at) + VALUES ('MatchChatMessage', 'Author', $3, 'user', $1::bigint, + 'match:m-1', + jsonb_build_object('messageId', $2::text), + $4, CASE WHEN $5 THEN now() END) + RETURNING id::text AS id`, + [ + steamId, + messageId, + overrides.message ?? "typo", + overrides.is_read ?? false, + overrides.deleted ?? false, + ], + ); + return inserted.id; + }; + + const text = async (id: string) => + ( + await postgres.query>( + `SELECT message FROM notifications WHERE id = $1::uuid`, + [id], + ) + )[0].message; + + it("shows the edited text on the recipient's unread row", async () => { + const user = await author(); + const reader = await fx.player("Reader"); + const matchId = randomUUID(); + const id = await place(matchId, user); + const unread = await row(reader, id); + + await edit(matchId, id, user, "fixed"); + + expect(await text(unread)).toBe("<b>fixed</b>"); + }); + + // A read row stays in the bell, and a collapsed one can be restored by + // its recipient, so neither may keep the text the author took back. + it("rewrites every row for that message, read or collapsed", async () => { + const reader = await fx.player("Reader"); + const messageId = randomUUID(); + const unread = await row(reader, messageId); + const read = await row(reader, messageId, { is_read: true }); + const collapsed = await row(reader, messageId, { deleted: true }); + const other = await row(reader, randomUUID()); + + await notifications.updateChatMessagePreview(messageId, "fixed"); + + expect(await text(unread)).toBe("fixed"); + expect(await text(read)).toBe("fixed"); + expect(await text(collapsed)).toBe("fixed"); + expect(await text(other)).toBe("typo"); + }); + + it("leaves a row retracted mid-edit blank", async () => { + const reader = await fx.player("Reader"); + const messageId = randomUUID(); + const id = await row(reader, messageId); + let pending: Promise = Promise.resolve(); + + await postgres.transaction(async (client) => { + await client.query( + `UPDATE notifications SET deleted_at = now(), message = '' + WHERE id = $1::uuid`, + [id], + ); + + pending = notifications.updateChatMessagePreview(messageId, "fixed"); + await new Promise((resolve) => setTimeout(resolve, 200)); + }); + + await pending; + + expect(await text(id)).toBe(""); + }); + + it("leaves a retracted row blank when the edit lands after the delete", async () => { + const reader = await fx.player("Reader"); + const messageId = randomUUID(); + const id = await row(reader, messageId); + + await notifications.retractChatMessage(messageId); + await notifications.updateChatMessagePreview(messageId, "fixed"); + + expect(await text(id)).toBe(""); + }); + + it("finds the message's rows through an index", async () => { + const plan = await postgres.transaction(async (client) => { + await client.query("SET LOCAL enable_seqscan = off"); + + const { rows } = await client.query( + `EXPLAIN UPDATE notifications SET message = 'x' + WHERE data->>'messageId' = $1 + AND type IN ('ChatMessage', 'MatchChatMessage') + AND message <> ''`, + [randomUUID()], + ); + + return rows.map((plan) => plan["QUERY PLAN"]).join("\n"); + }); + + expect(plan).toContain("notifications_message_id_idx"); + }); + }); +});