Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,8 @@ select_permissions:
- source
- deleted_by_steam_id
- deleted_at
- attachments
- gif
filter: {}
allow_aggregations: true
comment: Evidence of moderated website chat. Written only by the API.
Expand All @@ -38,6 +40,8 @@ select_permissions:
- source
- deleted_by_steam_id
- deleted_at
- attachments
- gif
filter:
room_type:
_neq: organizers
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ select_permissions:
name:
_nin:
- web_push_private_key
- giphy_api_key
comment: ""
- role: guest
permission:
Expand All @@ -33,7 +34,6 @@ update_permissions:
- role: administrator
permission:
columns:
- name
- value
filter: {}
check: {}
Expand Down
16 changes: 16 additions & 0 deletions hasura/migrations/default/1890000000200_chat_attachments/down.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
DROP TRIGGER IF EXISTS tad_direct_messages ON public.direct_messages;
DROP FUNCTION IF EXISTS public.tad_direct_messages();

-- The boot step only re-applies a triggers file whose digest changed, so
-- forget this one or a later up would never recreate it.
DELETE FROM migration_hashes.hashes
WHERE name = 'hasura/triggers/direct_messages';

ALTER TABLE public.direct_messages
DROP COLUMN IF EXISTS attachments,
DROP COLUMN IF EXISTS gif;

DROP TABLE IF EXISTS public.chat_attachments;

DELETE FROM public.settings
WHERE name IN ('chat_attachment_max_mb', 'giphy_api_key');
46 changes: 46 additions & 0 deletions hasura/migrations/default/1890000000200_chat_attachments/up.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
CREATE TABLE IF NOT EXISTS public.chat_attachments (
id uuid PRIMARY KEY DEFAULT gen_random_uuid(),
-- Kept when the player goes: the row is what the file is swept by.
uploader_steam_id bigint REFERENCES public.players (steam_id)
ON UPDATE CASCADE ON DELETE SET NULL,
room_type text NOT NULL,
room_id text NOT NULL,
storage_prefix text NOT NULL,
file_name text NOT NULL,
mime_type text NOT NULL,
size bigint NOT NULL CHECK (size > 0),
width integer,
height integer,
duration_ms integer,
poster_mime_type text,
-- Set until the multipart upload completes.
upload_id text,
uploaded_at timestamptz,
message_id uuid,
sent_at timestamptz,
-- NULL only for a sent direct message's file, which goes with its message.
expires_at timestamptz,
created_at timestamptz NOT NULL DEFAULT now()
);

CREATE INDEX IF NOT EXISTS chat_attachments_expires_at_idx
ON public.chat_attachments (expires_at)
WHERE expires_at IS NOT NULL;

CREATE INDEX IF NOT EXISTS chat_attachments_message_id_idx
ON public.chat_attachments (message_id)
WHERE message_id IS NOT NULL;

CREATE INDEX IF NOT EXISTS chat_attachments_pending_idx
ON public.chat_attachments (uploader_steam_id)
WHERE message_id IS NULL;

CREATE INDEX IF NOT EXISTS chat_attachments_room_idx
ON public.chat_attachments (room_type, room_id);

CREATE INDEX IF NOT EXISTS chat_attachments_storage_prefix_idx
ON public.chat_attachments (storage_prefix text_pattern_ops);

ALTER TABLE public.direct_messages
ADD COLUMN IF NOT EXISTS attachments jsonb,
ADD COLUMN IF NOT EXISTS gif jsonb;
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
ALTER TABLE public.chat_message_deletions
DROP COLUMN IF EXISTS attachments,
DROP COLUMN IF EXISTS gif;

DROP TABLE IF EXISTS public.chat_attachment_usage;

ALTER TABLE public.chat_attachments
DROP COLUMN IF EXISTS deleted_at;

DELETE FROM public.settings
WHERE name IN ('chat_attachment_daily_mb', 'giphy_hourly_limit');
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
-- Set when a group room's message is deleted: the file is kept as evidence,
-- for staff only, until it expires.
ALTER TABLE public.chat_attachments
ADD COLUMN IF NOT EXISTS deleted_at timestamptz;

-- What a player uploaded, kept apart from the files so removing a file does
-- not give its bytes back to the daily allowance.
CREATE TABLE IF NOT EXISTS public.chat_attachment_usage (
id bigserial PRIMARY KEY,
steam_id bigint NOT NULL REFERENCES public.players (steam_id)
ON UPDATE CASCADE ON DELETE CASCADE,
bytes bigint NOT NULL CHECK (bytes > 0),
created_at timestamptz NOT NULL DEFAULT now()
);

CREATE INDEX IF NOT EXISTS chat_attachment_usage_steam_id_created_at_idx
ON public.chat_attachment_usage (steam_id, created_at);

CREATE INDEX IF NOT EXISTS chat_attachment_usage_created_at_idx
ON public.chat_attachment_usage (created_at);

ALTER TABLE public.chat_message_deletions
ADD COLUMN IF NOT EXISTS attachments jsonb,
ADD COLUMN IF NOT EXISTS gif jsonb;
23 changes: 23 additions & 0 deletions hasura/triggers/direct_messages.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
-- Whatever deletes a direct message -- its author, the retention sweep, its
-- author's account going -- its files are due. The sweep removes them from
-- storage before it forgets the row.
CREATE OR REPLACE FUNCTION public.tad_direct_messages() RETURNS TRIGGER
LANGUAGE plpgsql
AS $$
BEGIN
UPDATE public.chat_attachments a
SET expires_at = now()
FROM deleted d
WHERE a.message_id = d.id
AND a.room_type = 'direct'
AND (a.expires_at IS NULL OR a.expires_at > now());

RETURN NULL;
END;
$$;

DROP TRIGGER IF EXISTS tad_direct_messages ON public.direct_messages;
CREATE TRIGGER tad_direct_messages
AFTER DELETE ON public.direct_messages
REFERENCING OLD TABLE AS deleted
FOR EACH STATEMENT EXECUTE FUNCTION public.tad_direct_messages();
154 changes: 154 additions & 0 deletions src/chat/chat-attachments-abort.spec.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,154 @@
import {
createServer,
request as httpRequest,
IncomingMessage,
Server,
} from "http";
import { AddressInfo } from "net";
import { Readable, Writable } from "stream";
import { ChatAttachmentsService } from "./chat-attachments.service";
import { ChatErrorCode } from "./enums/ChatErrorCode";

const ID = "0b7d6c1e-1111-4a2b-9c3d-000000000001";
const PART = ChatAttachmentsService.PART_SIZE;

// A real socket that goes away mid-body, against storage that does what the
// AWS SDK does with a body: pipe it, and listen for nothing else on it.
describe("ChatAttachmentsService, a part whose sender goes away", () => {
let server: Server;
let service: ChatAttachmentsService;
let sends: Array<{ signal?: AbortSignal; settled: Promise<unknown> }>;
let queries: string[];
let arrived: () => void;
let started: Promise<void>;
let result: Promise<ChatErrorCode | null>;

const row = {
id: ID,
uploader_steam_id: "1",
room_type: "matchmaking",
room_id: "lobby-1",
storage_prefix: `chat-attachments/rooms/2026-10-02/${ID}/`,
file_name: "clip.mp4",
mime_type: "video/mp4",
size: String(PART * 3),
width: null,
height: null,
duration_ms: null,
poster_mime_type: null,
upload_id: "upload-1",
message_id: null,
deleted_at: null,
};

const s3 = {
uploadPart: jest.fn(
(
_key: string,
_uploadId: string,
_part: number,
body: Readable,
_length: number,
signal?: AbortSignal,
) => {
const settled = new Promise<void>((resolve, reject) => {
const sink = new Writable({
write(_chunk, _encoding, callback) {
arrived();
callback();
},
});

body.pipe(sink);
sink.on("finish", () => resolve());
signal?.addEventListener("abort", () =>
reject(new Error("request aborted")),
);
});

sends.push({ signal, settled: settled.catch(() => {}) });

return settled;
},
),
};

const postgres = {
query: jest.fn(async (sql: string) => {
queries.push(sql);
return sql.includes("upload_id IS NOT NULL") ? [row] : [];
}),
};

beforeEach(async () => {
sends = [];
queries = [];
started = new Promise((resolve) => {
arrived = resolve;
});

service = new ChatAttachmentsService(
{ log: jest.fn(), warn: jest.fn(), error: jest.fn() } as any,
postgres as any,
s3 as any,
);

server = createServer((request: IncomingMessage) => {
result = service.uploadPart("1", ID, 2, request, PART);
});

await new Promise<void>((resolve) => server.listen(0, resolve));
});

afterEach(async () => {
await new Promise((resolve) => server.close(resolve));
});

const abortMidBody = async () => {
const { port } = server.address() as AddressInfo;

const client = httpRequest({
port,
method: "PUT",
headers: {
"content-type": "application/octet-stream",
"content-length": String(PART),
},
});
client.on("error", () => {});

client.write(Buffer.alloc(300 * 1024));
await started;
client.destroy();

return await result;
};

it("survives, refuses the part, and lets go of the request to storage", async () => {
await expect(abortMidBody()).resolves.toBe(ChatErrorCode.Invalid);

expect(sends).toHaveLength(1);
expect(sends[0].signal?.aborted).toBe(true);
await sends[0].settled;
});

it("frees the player's upload slots for the next try", async () => {
await abortMidBody();

expect((service as any).uploadsInFlight).toBe(0);
expect((service as any).uploadsByPlayer.size).toBe(0);
});

it("leaves the upload for the sweep to abort", async () => {
await abortMidBody();

expect(
queries.filter(
(sql) =>
sql.includes("DELETE FROM public.chat_attachments") ||
sql.includes("SET expires_at") ||
sql.includes("SET upload_id = NULL"),
),
).toEqual([]);
});
});
Loading
Loading