diff --git a/README.md b/README.md index 280e57f..797042c 100644 --- a/README.md +++ b/README.md @@ -16,7 +16,7 @@ - A janela do MVP exibe progresso, estados de sucesso/erro e pode ser fechada pelo botão visível ou pela tecla `Esc`. No Linux/Wayland, a Waybar abre a janela flutuante compacta em cerca de `432x272` no compositor (`380x220` de área interna Tauri) e fechar encerra a janela como antes; no Windows, a janela usa a mesma área interna compacta, fechar oculta a janela, mantém o app vivo na tray até o usuário escolher `Sair`, abre posicionada acima da área da tray e pode ser arrastada pela barra superior customizada. No primeiro start no Windows, o app ativa `Iniciar com Windows` automaticamente e grava um marcador local; se o usuário desativar o autostart no menu da tray, o app não reativa sozinho em starts futuros. Fonte: launcher Waybar `scripts/quickdrop-waybar`, configuração Tauri `src-tauri/tauri.conf.json`, tray/posicionamento/autostart em `src-tauri/src/lib.rs` e UI `src/desktop/App.tsx`. - A instalação Windows por PowerShell é pública no endpoint `GET /install.ps1`; o `.exe` é baixado pelo endpoint interno `GET /windows/latest.exe`, que usa um token GitHub configurado somente no servidor para buscar o asset privado `QuickDrop_*_x64-setup.exe` da última release sem expor credenciais ao usuário final. A página principal exibe `irm https://quickdrop.eaedave.xyz/install.ps1 | iex` com botão de cópia. Após o NSIS silencioso concluir, o script abre o app instalado em modo visível no canto direito e libera o terminal. Fonte: `src/desktop/App.tsx`, `scripts/install-windows.ps1`, `src/server/index.ts` e `src/server/windows-installer-service.ts`. - A instalação Linux por Bash é pública no endpoint `GET /install.sh` (`curl -fsSL https://quickdrop.eaedave.xyz/install.sh | bash`); o script detecta Linux/x86_64 com Waybar (e avisa se faltar Hyprland ou dependências de runtime como webkit2gtk, gtk3, wl-clipboard e libnotify), baixa o binário pré-compilado pelo endpoint interno `GET /linux/latest` — que usa o mesmo token GitHub server-side para buscar o asset privado `quickdrop_*_x86_64-linux` da última release —, instala `~/.local/bin/quickdrop` e o launcher `~/.local/bin/quickdrop-waybar` (servido por `GET /linux/quickdrop-waybar`), e registra de forma idempotente o módulo `custom/quickdrop` na Waybar com backup do config. Fonte: `scripts/install-linux.sh`, `src/server/index.ts`, `src/server/linux-installer-service.ts` e `src/server/github-release.ts`. -- O QuickDrop tem um relay de texto em tempo real para colar/compartilhar texto entre máquinas sem clipboard compartilhado (ex.: máquinas Guacamole). Pela web, o usuário cria uma sala (botão "Criar nova sala") ou entra com um código curto de 6 caracteres; um textarea grande é sincronizado ao vivo entre todos na mesma sala via WebSocket, no modelo último-a-escrever-vence (last-writer-wins) com versão monotônica. Edições simultâneas não sobrescrevem em silêncio: quando chega uma alteração remota durante uma edição local pendente, a UI mostra um aviso não destrutivo com opção de carregar. As salas ficam em memória (não sobrevivem a redeploy), expiram após inatividade (padrão 12h sem clientes) e limitam tamanho do texto (padrão 256 KB) e número de salas/clientes. A página é servida pelo mesmo backend (subdomínio `texto.*`, rota `/t` ou `?c=CÓDIGO` na raiz). Fontes: Endpoints internos `POST /api/text` (cria sala), `GET /api/text/:code` (snapshot) e `WS /api/text/:code/ws` (sync); redirects internos `GET /t` e `GET /t/:code`; store `src/server/text-session-store.ts`, serviço `src/server/text-session-service.ts`, UI `src/desktop/TextSession.tsx`. +- O QuickDrop tem um relay de texto em tempo real para colar/compartilhar texto entre máquinas sem clipboard compartilhado (ex.: máquinas Guacamole). Pela web, o usuário cria uma sala (botão "Criar nova sala") ou entra com um código curto de 6 caracteres; um textarea grande é sincronizado ao vivo entre todos na mesma sala via WebSocket, no modelo último-a-escrever-vence (last-writer-wins) com versão monotônica. Edições simultâneas não sobrescrevem em silêncio: quando chega uma alteração remota durante uma edição local pendente, a UI mostra um aviso não destrutivo com opção de carregar. As salas agora ficam persistidas no PostgreSQL, então sobrevivem a restart/deploy; continuam expirando por inatividade (padrão 12h sem clientes), limitam tamanho do texto (padrão 256 KB) e número de salas/clientes. A página é servida pelo mesmo backend (subdomínio `texto.*`, rota `/t` ou `?c=CÓDIGO` na raiz). Fontes: Endpoints internos `POST /api/text` (cria sala), `GET /api/text/:code` (snapshot) e `WS /api/text/:code/ws` (sync); redirects internos `GET /t` e `GET /t/:code`; repositório `src/server/text-rooms-repository.ts`, hub `src/server/text-session-hub.ts`, serviço `src/server/text-session-service.ts`, UI `src/desktop/TextSession.tsx`. @@ -99,7 +99,7 @@ Relay de texto em tempo real entre máquinas, no mesmo backend: - `GET /api/text/:code` — snapshot atual `{ text, version }` (404 se a sala não existe). - `WS /api/text/:code/ws` — o servidor envia `snapshot` ao conectar, `update` quando outro cliente escreve, `ack` ao autor após cada escrita e `error` (texto acima do limite ou sala inválida). O cliente envia `{ type: "write", text, baseVersion }`. Heartbeat ping/pong derruba conexões mortas. -Estado em memória (perde no redeploy), com varredura de TTL e limites configuráveis por ambiente: +Estado persistido no PostgreSQL (sobrevive a restart/redeploy), com expiração por TTL e limites configuráveis por ambiente: ```env TEXT_SESSION_TTL_HOURS=12 diff --git a/docs/LLM_CONTEXT.md b/docs/LLM_CONTEXT.md index 9d76346..dc39fca 100644 --- a/docs/LLM_CONTEXT.md +++ b/docs/LLM_CONTEXT.md @@ -13,7 +13,7 @@ - Regra: no Windows, QuickDrop roda residente na system tray. Clique esquerdo ou menu `Abrir QuickDrop` abre/foca a janela; botão fechar/Esc ocultam a janela; menu `Sair` encerra o processo; o primeiro start ativa autostart automaticamente com `--tray-start` e grava marcador local para não reativar sozinho se o usuário desativar depois via menu `Iniciar com Windows`. Fonte: `src-tauri/src/lib.rs`, `src-tauri/tauri.windows.conf.json` e `src/desktop/App.tsx`. - Regra: a instalação Windows por PowerShell (`GET /install.ps1`) baixa o instalador por `GET /windows/latest.exe`; esse endpoint usa `QUICKDROP_GITHUB_TOKEN`/`GITHUB_TOKEN` somente no backend para buscar o asset privado `QuickDrop_*_x64-setup.exe` da última release GitHub, sem expor token ao script público. `GET /` mostra o mesmo comando para cópia rápida. Depois do NSIS silencioso, o script abre o `QuickDrop.exe` instalado sem `--tray-start`, deixando a janela visível no canto direito e liberando o terminal sem esperar o processo residente da tray encerrar. Fonte: `src/desktop/App.tsx`, `scripts/install-windows.ps1`, `src/server/index.ts`, `src/server/windows-installer-service.ts`, `src/server/index.test.ts`. - Regra: a instalação Linux por Bash (`GET /install.sh`, `curl -fsSL https://quickdrop.eaedave.xyz/install.sh | bash`) detecta Linux x86_64 com Waybar (avisa se faltar Hyprland ou as dependências de runtime webkit2gtk/gtk3/wl-clipboard/libnotify), baixa o binário pré-compilado por `GET /linux/latest` (mesmo token GitHub server-side, asset privado `quickdrop_*_x86_64-linux`), instala `~/.local/bin/quickdrop` + launcher `~/.local/bin/quickdrop-waybar` (de `GET /linux/quickdrop-waybar`) e faz patch idempotente do módulo `custom/quickdrop` em `~/.config/waybar/config.jsonc` com backup e restart. Fonte: `scripts/install-linux.sh`, `src/server/index.ts`, `src/server/linux-installer-service.ts`, `src/server/github-release.ts`, `src/server/index.test.ts`, `scripts/install-linux.test.ts`. -- Regra: QuickDrop também oferece um relay de texto em tempo real para colar/copiar texto entre máquinas sem clipboard compartilhado. A web cria sala por `POST /api/text`, compartilha um código curto de 6 caracteres e sincroniza um textarea grande via `WS /api/text/:code/ws`; o snapshot inicial vem de `GET /api/text/:code`. O modelo é last-writer-wins com `version` monotônica, aviso não destrutivo quando chega atualização remota durante edição local pendente, limite padrão de 256 KB por sala, expiração após 12 horas sem clientes e armazenamento só em memória (perde no redeploy). A entrada pode acontecer por host `texto.*`, por `?c=CÓDIGO` na raiz ou por `GET /t`/`GET /t/:code`, que redirecionam para a SPA com `?c=`. Fonte: endpoints internos `/api/text`, `/api/text/:code`, `WS /api/text/:code/ws`, redirects `/t` e `/t/:code`; `src/server/text-session-store.ts`; `src/server/text-session-service.ts`; `src/desktop/TextSession.tsx`. +- Regra: QuickDrop também oferece um relay de texto em tempo real para colar/copiar texto entre máquinas sem clipboard compartilhado. A web cria sala por `POST /api/text`, compartilha um código curto de 6 caracteres e sincroniza um textarea grande via `WS /api/text/:code/ws`; o snapshot inicial vem de `GET /api/text/:code`. O modelo é last-writer-wins com `version` monotônica, aviso não destrutivo quando chega atualização remota durante edição local pendente, limite padrão de 256 KB por sala e expiração após 12 horas sem clientes. Diferente do MVP inicial, o estado da sala agora fica persistido em PostgreSQL e sobrevive a restart/deploy; o processo só mantém o hub de conexões WebSocket em memória. A entrada pode acontecer por host `texto.*`, por `?c=CÓDIGO` na raiz ou por `GET /t`/`GET /t/:code`, que redirecionam para a SPA com `?c=`. Fonte: endpoints internos `/api/text`, `/api/text/:code`, `WS /api/text/:code/ws`, redirects `/t` e `/t/:code`; `src/server/text-rooms-repository.ts`; `src/server/text-session-hub.ts`; `src/server/text-session-service.ts`; `src/desktop/TextSession.tsx`. ## Technical map for future LLMs @@ -25,7 +25,7 @@ - Download workflow: `src/server/download-service.ts` rejects missing/deleted, deletes+marks expired, otherwise signs R2 GET and redirects 302. - Desktop Rust commands: `src-tauri/src/lib.rs` implements `upload_file`, `upload_files`, `read_clipboard_upload_inputs`, `copy_link`, `notify_success`, `dismiss_window`, `uses_native_clipboard_paste`; `upload_files` validates selected paths, creates a temporary `.zip` with unique entry names when there is more than 1 file, uploads that ZIP through the existing single-file backend endpoint, then removes generated temp files. `read_clipboard_upload_inputs` uses `wl-paste --list-types`/`--type` on Linux/Wayland to materialize clipboard image/text as a temporary local file for upload. Config reads `QUICKDROP_API_BASE_URL`; default is `https://quickdrop.eaedave.xyz` on Windows and `http://127.0.0.1:3000` on non-Windows unless the launcher/env overrides it. - Desktop UI / Web Frontend: `src/desktop/App.tsx` funciona em ambos os ambientes (detectado via `isTauri`). No desktop usa drag-drop events do Tauri, seletor nativo e `Ctrl+V` via `wl-paste` somente quando `uses_native_clipboard_paste` retorna true; na web renderiza uma landing page com hero, card de instalação com abas Windows/Linux (comando + link `Ver script` para `/install.ps1` ou `/install.sh` conforme a aba selecionada via `INSTALL_PLATFORMS`), upload HTML5 por drag-and-drop/input e `Ctrl+V` via Clipboard API/eventos padrão. Mostra progresso real do upload e fase `Criando ZIP...` ao compactar múltiplos arquivos. A barra superior customizada chama `getCurrentWindow().startDragging()` para mover a janela frameless no Tauri. -- Text relay web-only: `src/desktop/main.tsx` desvia para `src/desktop/TextSession.tsx` quando `src/desktop/web-route.ts` detecta host `texto.*`, rota `/t`/`/t/:code` ou `?c=CÓDIGO`; o restante continua em `App.tsx`. `src/desktop/text-client.ts` fala com `POST /api/text`, `GET /api/text/:code` e `WS /api/text/:code/ws`, implementa reconnect com backoff e normaliza payloads sem dependências extras. `TextSession.tsx` mantém textarea local + debounce de ~150 ms, suprime eco próprio por `clientId`, mostra banner de conflito remoto e alterna a classe `quickdrop-text-page` no ``. O backend usa `@fastify/websocket` em `src/server/index.ts`, `TextSessionStore` em memória e os helpers `startTextSessionSweep`/`startTextSessionHeartbeat`; `src/server/index.ts` também redireciona `/t` e `/t/:code` para `/?c=` para a SPA carregar sem fallback especial do `@fastify/static`. Envs novos: `TEXT_SESSION_TTL_HOURS`, `TEXT_SESSION_MAX_KB`, `TEXT_SESSION_CODE_LENGTH`, `TEXT_SESSION_MAX_SESSIONS`, `TEXT_SESSION_MAX_CLIENTS`. +- Text relay web-only: `src/desktop/main.tsx` desvia para `src/desktop/TextSession.tsx` quando `src/desktop/web-route.ts` detecta host `texto.*`, rota `/t`/`/t/:code` ou `?c=CÓDIGO`; o restante continua em `App.tsx`. `src/desktop/text-client.ts` fala com `POST /api/text`, `GET /api/text/:code` e `WS /api/text/:code/ws`, implementa reconnect com backoff e normaliza payloads sem dependências extras. `TextSession.tsx` mantém textarea local + debounce de ~150 ms, suprime eco próprio por `clientId`, mostra banner de conflito remoto e alterna a classe `quickdrop-text-page` no ``. O backend agora separa `src/server/text-session-hub.ts` (clientes WS em memória) de `src/server/text-rooms-repository.ts` (estado persistido em Postgres via Drizzle), rearma salas abertas no restart com `rearmTextRoomsAfterRestart`, e usa `startTextSessionSweep` para marcar expiradas em DB; `src/server/index.ts` continua redirecionando `/t` e `/t/:code` para `/?c=` para a SPA carregar sem fallback especial do `@fastify/static`. Envs: `TEXT_SESSION_TTL_HOURS`, `TEXT_SESSION_MAX_KB`, `TEXT_SESSION_CODE_LENGTH`, `TEXT_SESSION_MAX_SESSIONS`, `TEXT_SESSION_MAX_CLIENTS`. - Commands: `bun run server:dev`, `bun run server:start`, `bun run db:generate`, `bun run db:check`, `bun run db:migrate`, `bun run cleanup:run`, `bun run desktop:dev`, `bun run desktop:build:web`, `bun run desktop:build`, `bun run desktop:build:windows`, `bun run desktop:package:linux`, `bun run desktop:install`, `bun run desktop:install:local`, `bun run waybar:install`, `bun run local:server:up`, `bun run local:server:stop`, `bun run quickdrop:install`, `bun run quickdrop:install:local`, `bun run typecheck`, `bun test`. - Waybar install: `scripts/install-desktop.ts` copies `quickdrop`; `scripts/install-waybar-module.ts` installs/refreshes `quickdrop-waybar`, idempotently patches `~/.config/waybar/config.jsonc` with `QUICKDROP_API_BASE_URL=https://quickdrop.eaedave.xyz` by default, backs it up, and restarts Waybar via `omarchy restart waybar` when available. The Tauri window uses compact `380x220` inner size; the Waybar launcher positions it as about `432x272` compositor size by default, supports `QUICKDROP_WINDOW_WIDTH`/`QUICKDROP_WINDOW_HEIGHT` overrides, and re-places the window after launch so Tauri's own initial position cannot leave the top under the bar. Override backend with `QUICKDROP_API_BASE_URL=... bun run waybar:install`; `quickdrop:install` builds+installs the desktop against the remote backend for reboot-safe daily use. - Windows tray: `src-tauri/Cargo.toml` enables Tauri `tray-icon`/`image-png` and native plugins for clipboard, notification, and autostart. `src-tauri/src/lib.rs` builds the tray only on Windows, uses `icons/icon.png`, positions the default/opened window above the tray/work-area corner instead of at the cursor, hides instead of closes on `dismiss_window`, enables autostart once on first launch, persists `autostart-configured` under the app config dir, and autostarts with `--tray-start`. `scripts/install-windows.ps1` intentionally does not pass `--tray-start` after installation: it waits for NSIS only, then starts the installed `QuickDrop.exe` detached/visible so the PowerShell prompt returns immediately. `src-tauri/tauri.windows.conf.json` enables NSIS bundling and Windows icons. @@ -60,4 +60,5 @@ - 2026-06-23: Waybar ficou mais compacta de verdade: o binário Tauri agora usa área interna `380x220`, o launcher posiciona como ~`432x272` e reaplica a posição após a janela existir para não deixar o topo sob a barra; `bun run waybar:install` também instala/atualiza `~/.local/bin/quickdrop-waybar`, não só o JSONC da Waybar. Fontes atualizadas: `src-tauri/src/lib.rs`, `scripts/quickdrop-waybar`, `scripts/install-waybar-module.ts`, `scripts/install-desktop.ts`, `scripts/install-waybar-module.test.ts`, README e LLM context. - 2026-06-23: Adicionada instalação Linux (Hyprland + Waybar) espelhando o fluxo Windows: novo `scripts/install-linux.sh` servido por `GET /install.sh` (`curl -fsSL https://quickdrop.eaedave.xyz/install.sh | bash`) detecta ambiente, baixa o binário por `GET /linux/latest` e instala `~/.local/bin/quickdrop` + launcher (de `GET /linux/quickdrop-waybar`) com patch idempotente da Waybar via patcher Python embarcado (paridade testada com `install-waybar-module.ts`). Extraído helper compartilhado `src/server/github-release.ts` (usado por `windows-installer-service.ts` e novo `linux-installer-service.ts`); landing com abas Windows/Linux funcionais; `scripts/package-linux-release.ts` + `bun run desktop:package:linux` geram `quickdrop__x86_64-linux`, publicado na release `v0.1.1` (latest). Fontes: `src/server/index.ts`, `src/server/index.test.ts`, `scripts/install-linux.sh`, `scripts/install-linux.test.ts`, `scripts/package-linux-release.ts`, `src/desktop/App.tsx`, `src/desktop/input.css`, `Dockerfile`, `package.json`, README e LLM context. - 2026-06-23: Adicionado relay de texto em tempo real no mesmo backend Fastify usando `@fastify/websocket`, com salas efêmeras em memória (`POST /api/text`, `GET /api/text/:code`, `WS /api/text/:code/ws`), heartbeat/sweep configuráveis por env e testes WS reais em `src/server/text-session-service.test.ts`. No frontend web, `main.tsx` agora roteia para `TextSession.tsx` via `web-route.ts` quando o host é `texto.*`, a rota é `/t`/`/t/:code` ou há `?c=`, e `text-client.ts` faz reconnect/backoff e sincronização do textarea com banner de conflito remoto. Na sequência, `src/server/index.ts` ganhou redirects `GET /t` e `GET /t/:code` para `/?c=` porque o backend serve a SPA de `dist/` e a rota path-based precisava de entrada server-side explícita. Fontes alteradas: `src/server/index.ts`, `src/server/index.test.ts`, `src/server/config.ts`, `src/server/ids.ts`, `src/server/text-session-store.ts`, `src/server/text-session-service.ts`, `src/server/text-session-store.test.ts`, `src/server/text-session-service.test.ts`, `src/server/r2.test.ts`, `src/desktop/main.tsx`, `src/desktop/web-route.ts`, `src/desktop/text-client.ts`, `src/desktop/TextSession.tsx`, `src/desktop/input.css`, `.env.example`, README e LLM context. +- 2026-06-23: PR 1 da fase 2 do relay de texto: o estado das salas saiu do `Map` em memória e passou para PostgreSQL/Drizzle com tabela `text_rooms` (`migrations/002_create_text_rooms.sql`, idempotente, já que `migrate.ts` reaplica todos os `.sql` em cada start). O backend agora separa `TextSessionHub` (clientes WS) de `textRoomsRepository` (código/texto/version/expiração), rearma salas `expires_at = null` após restart e mantém o contrato HTTP/WS do MVP. Fontes alteradas: `src/server/schema.ts`, `src/server/text-rooms-repository.ts`, `src/server/text-session-hub.ts`, `src/server/text-session-service.ts`, `src/server/index.ts`, `migrations/002_create_text_rooms.sql`, testes do relay, README e LLM context. diff --git a/migrations/002_create_text_rooms.sql b/migrations/002_create_text_rooms.sql new file mode 100644 index 0000000..316c512 --- /dev/null +++ b/migrations/002_create_text_rooms.sql @@ -0,0 +1,14 @@ +create table if not exists text_rooms ( + code varchar(16) primary key, + text text not null default '', + version integer not null default 0, + pin_hash text, + created_at timestamptz not null, + updated_at timestamptz not null, + expires_at timestamptz, + deleted_at timestamptz +); + +create index if not exists text_rooms_expires_idx + on text_rooms (expires_at) + where deleted_at is null; diff --git a/src/server/index.ts b/src/server/index.ts index 54beded..461f6c2 100644 --- a/src/server/index.ts +++ b/src/server/index.ts @@ -13,8 +13,13 @@ import { createR2Client } from "./r2"; import { handleUpload } from "./upload-service"; import { handleLinuxInstallerDownload } from "./linux-installer-service"; import { handleWindowsInstallerDownload } from "./windows-installer-service"; -import { TextSessionStore } from "./text-session-store"; -import { registerTextSessionRoutes, startTextSessionHeartbeat, startTextSessionSweep } from "./text-session-service"; +import { TextSessionHub } from "./text-session-hub"; +import { + rearmTextRoomsAfterRestart, + registerTextSessionRoutes, + startTextSessionHeartbeat, + startTextSessionSweep, +} from "./text-session-service"; export function buildApp() { const config = loadConfig(); @@ -36,11 +41,8 @@ export function buildApp() { options: { maxPayload: config.textSessionMaxBytes + 1024 }, }); - const textStore = new TextSessionStore({ - maxBytes: config.textSessionMaxBytes, - maxSessions: config.textSessionMaxSessions, + const textHub = new TextSessionHub({ maxClientsPerSession: config.textSessionMaxClientsPerSession, - codeLength: config.textSessionCodeLength, }); const sendScriptFile = (reply: FastifyReply, fileName: string) => @@ -85,21 +87,28 @@ export function buildApp() { async (request, reply) => handleUpload(request, reply, { config, r2Client }), ); - registerTextSessionRoutes(app, textStore); + registerTextSessionRoutes(app, { + hub: textHub, + maxBytes: config.textSessionMaxBytes, + maxSessions: config.textSessionMaxSessions, + codeLength: config.textSessionCodeLength, + ttlMs: config.textSessionTtlHours * 60 * 60 * 1000, + }); }); app.get<{ Params: { shortId: string } }>("/f/:shortId", async (request, reply) => { await handleDownload(request.params.shortId, reply, { config, r2Client }); }); - return { app, config, textStore }; + return { app, config, textHub }; } export async function startServer(): Promise { - const { app, config, textStore } = buildApp(); + const { app, config } = buildApp(); + await rearmTextRoomsAfterRestart(config.textSessionTtlHours * 60 * 60 * 1000); await app.listen({ host: "0.0.0.0", port: config.port }); startCleanupJob(); - startTextSessionSweep(textStore, config.textSessionTtlHours * 60 * 60 * 1000); + startTextSessionSweep(); startTextSessionHeartbeat(app); } diff --git a/src/server/schema.ts b/src/server/schema.ts index 421d215..6dc91ee 100644 --- a/src/server/schema.ts +++ b/src/server/schema.ts @@ -21,3 +21,22 @@ export const uploads = pgTable( ); export type UploadRecord = typeof uploads.$inferSelect; + +export const textRooms = pgTable( + "text_rooms", + { + code: varchar("code", { length: 16 }).primaryKey(), + text: text("text").notNull().default(""), + version: integer("version").default(0).notNull(), + pinHash: text("pin_hash"), + createdAt: timestamp("created_at", { withTimezone: true }).notNull(), + updatedAt: timestamp("updated_at", { withTimezone: true }).notNull(), + expiresAt: timestamp("expires_at", { withTimezone: true }), + deletedAt: timestamp("deleted_at", { withTimezone: true }), + }, + (table) => [ + index("text_rooms_expires_idx").on(table.expiresAt).where(sql`${table.deletedAt} is null`), + ], +); + +export type TextRoomRecord = typeof textRooms.$inferSelect; diff --git a/src/server/text-rooms-repository.ts b/src/server/text-rooms-repository.ts new file mode 100644 index 0000000..452470e --- /dev/null +++ b/src/server/text-rooms-repository.ts @@ -0,0 +1,187 @@ +import { and, asc, eq, gt, isNull, lte, or, sql as drizzleSql } from "drizzle-orm"; +import { db } from "./db"; +import { textRooms, type TextRoomRecord } from "./schema"; + +export type TextRoomRow = { + code: string; + text: string; + version: number; + pin_hash: string | null; + created_at: Date; + updated_at: Date; + expires_at: Date | null; + deleted_at: Date | null; +}; + +export type TextRoomsRepository = { + countActiveTextRooms(now: Date): Promise; + createTextRoom(input: { + code: string; + text: string; + version: number; + pinHash?: string | null; + createdAt: Date; + updatedAt: Date; + expiresAt: Date; + }): Promise; + findTextRoomByCode(code: string): Promise; + markTextRoomActive(code: string, updatedAt: Date): Promise; + updateTextRoomText(input: { code: string; text: string; now: Date }): Promise; + scheduleTextRoomExpiry(code: string, expiresAt: Date, updatedAt: Date): Promise; + rearmOpenTextRooms(expiresAt: Date, updatedAt: Date): Promise; + findExpiredTextRooms(now: Date, limit?: number): Promise; + markTextRoomDeleted(code: string, deletedAt: Date): Promise; +}; + +export async function countActiveTextRooms(now: Date): Promise { + const rows = await db + .select({ count: drizzleSql`count(*)` }) + .from(textRooms) + .where(activeRoomFilter(now)); + const row = rows[0]; + + return Number(row?.count ?? 0); +} + +export async function createTextRoom(input: { + code: string; + text: string; + version: number; + pinHash?: string | null; + createdAt: Date; + updatedAt: Date; + expiresAt: Date; +}): Promise { + try { + const rows = await db + .insert(textRooms) + .values({ + code: input.code, + text: input.text, + version: input.version, + pinHash: input.pinHash ?? null, + createdAt: input.createdAt, + updatedAt: input.updatedAt, + expiresAt: input.expiresAt, + }) + .returning(); + const row = rows[0]; + + if (!row) { + throw new Error("Insert did not return a text room row"); + } + + return toTextRoomRow(row); + } catch (error) { + if (hasPostgresUniqueViolation(error)) { + return null; + } + + throw error; + } +} + +export async function findTextRoomByCode(code: string): Promise { + const rows = await db + .select() + .from(textRooms) + .where(and(eq(textRooms.code, code), isNull(textRooms.deletedAt))) + .limit(1); + const row = rows[0]; + + return row ? toTextRoomRow(row) : null; +} + +export async function markTextRoomActive(code: string, updatedAt: Date): Promise { + await db + .update(textRooms) + .set({ updatedAt, expiresAt: null }) + .where(and(eq(textRooms.code, code), isNull(textRooms.deletedAt))); +} + +export async function updateTextRoomText(input: { + code: string; + text: string; + now: Date; +}): Promise { + const rows = await db + .update(textRooms) + .set({ + text: input.text, + updatedAt: input.now, + version: drizzleSql`${textRooms.version} + 1`, + }) + .where(and(eq(textRooms.code, input.code), isNull(textRooms.deletedAt))) + .returning(); + const row = rows[0]; + + return row ? toTextRoomRow(row) : null; +} + +export async function scheduleTextRoomExpiry(code: string, expiresAt: Date, updatedAt: Date): Promise { + await db + .update(textRooms) + .set({ expiresAt, updatedAt }) + .where(and(eq(textRooms.code, code), isNull(textRooms.deletedAt))); +} + +export async function rearmOpenTextRooms(expiresAt: Date, updatedAt: Date): Promise { + await db + .update(textRooms) + .set({ expiresAt, updatedAt }) + .where(and(isNull(textRooms.deletedAt), isNull(textRooms.expiresAt))); +} + +export async function findExpiredTextRooms(now: Date, limit = 100): Promise { + const rows = await db + .select() + .from(textRooms) + .where(and(isNull(textRooms.deletedAt), lte(textRooms.expiresAt, now))) + .orderBy(asc(textRooms.expiresAt)) + .limit(limit); + + return rows.map(toTextRoomRow); +} + +export async function markTextRoomDeleted(code: string, deletedAt: Date): Promise { + await db + .update(textRooms) + .set({ deletedAt }) + .where(and(eq(textRooms.code, code), isNull(textRooms.deletedAt))); +} + +function activeRoomFilter(now: Date) { + return and( + isNull(textRooms.deletedAt), + or(isNull(textRooms.expiresAt), gt(textRooms.expiresAt, now)), + ); +} + +function toTextRoomRow(row: TextRoomRecord): TextRoomRow { + return { + code: row.code, + text: row.text, + version: row.version, + pin_hash: row.pinHash, + created_at: row.createdAt, + updated_at: row.updatedAt, + expires_at: row.expiresAt, + deleted_at: row.deletedAt, + }; +} + +function hasPostgresUniqueViolation(error: unknown): boolean { + return !!error && typeof error === "object" && "code" in error && error.code === "23505"; +} + +export const textRoomsRepository: TextRoomsRepository = { + countActiveTextRooms, + createTextRoom, + findTextRoomByCode, + markTextRoomActive, + updateTextRoomText, + scheduleTextRoomExpiry, + rearmOpenTextRooms, + findExpiredTextRooms, + markTextRoomDeleted, +}; diff --git a/src/server/text-session-hub.ts b/src/server/text-session-hub.ts new file mode 100644 index 0000000..3209cb6 --- /dev/null +++ b/src/server/text-session-hub.ts @@ -0,0 +1,82 @@ +import { normalizeSessionCode } from "./ids"; + +export type SessionClient = { + id: string; + send: (data: string) => void; + close: (code?: number, reason?: string) => void; +}; + +type HubRoom = { + code: string; + clients: Map; +}; + +export type TextSessionHubOptions = { + maxClientsPerSession: number; +}; + +export type JoinResult = + | { ok: true; clientCount: number } + | { ok: false; reason: "full" }; + +export class TextSessionHub { + private readonly rooms = new Map(); + private readonly maxClientsPerSession: number; + + constructor(options: TextSessionHubOptions) { + this.maxClientsPerSession = options.maxClientsPerSession; + } + + join(code: string, client: SessionClient): JoinResult { + const normalized = normalizeSessionCode(code); + const room = this.rooms.get(normalized) ?? { code: normalized, clients: new Map() }; + + if (room.clients.size >= this.maxClientsPerSession) { + return { ok: false, reason: "full" }; + } + + room.clients.set(client.id, client); + this.rooms.set(normalized, room); + return { ok: true, clientCount: room.clients.size }; + } + + leave(code: string, clientId: string): number { + const room = this.rooms.get(normalizeSessionCode(code)); + + if (!room) { + return 0; + } + + room.clients.delete(clientId); + if (room.clients.size === 0) { + this.rooms.delete(room.code); + return 0; + } + + return room.clients.size; + } + + clientCount(code: string): number { + return this.rooms.get(normalizeSessionCode(code))?.clients.size ?? 0; + } + + broadcast(code: string, data: string, exceptId?: string): void { + const room = this.rooms.get(normalizeSessionCode(code)); + + if (!room) { + return; + } + + for (const client of room.clients.values()) { + if (client.id === exceptId) { + continue; + } + + try { + client.send(data); + } catch { + // A failing client must not break delivery to the rest of the room. + } + } + } +} diff --git a/src/server/text-session-service.test.ts b/src/server/text-session-service.test.ts index c356269..a8c6b90 100644 --- a/src/server/text-session-service.test.ts +++ b/src/server/text-session-service.test.ts @@ -1,16 +1,14 @@ -import { afterAll, beforeAll, describe, expect, test } from "bun:test"; -import type { FastifyInstance } from "fastify"; -import { buildApp } from "./index"; - -const testEnv = { - PORT: "3000", - DATABASE_URL: "postgres://quickdrop:quickdrop@127.0.0.1:5432/quickdrop", - R2_ACCOUNT_ID: "account", - R2_ACCESS_KEY_ID: "access-key", - R2_SECRET_ACCESS_KEY: "secret-key", - PUBLIC_BASE_URL: "https://quickdrop.eaedave.xyz", - TEXT_SESSION_MAX_KB: "1", -}; +import { afterEach, describe, expect, test } from "bun:test"; +import Fastify, { type FastifyInstance } from "fastify"; +import rateLimit from "@fastify/rate-limit"; +import websocket from "@fastify/websocket"; +import { TextSessionHub } from "./text-session-hub"; +import { + rearmTextRoomsAfterRestart, + registerTextSessionRoutes, + startTextSessionSweep, +} from "./text-session-service"; +import type { TextRoomRow, TextRoomsRepository } from "./text-rooms-repository"; type ServerMessage = { type: "snapshot" | "update" | "ack" | "error"; @@ -21,45 +19,174 @@ type ServerMessage = { message?: string; }; -let app: FastifyInstance; -let wsBase: string; -const openSockets = new Set(); -const previousEnv: Record = {}; +type MutableClock = { current: Date }; + +class InMemoryTextRoomsRepository implements TextRoomsRepository { + private readonly rooms = new Map(); + + async countActiveTextRooms(now: Date): Promise { + let count = 0; + + for (const room of this.rooms.values()) { + if (room.deleted_at === null && (room.expires_at === null || room.expires_at > now)) { + count += 1; + } + } -beforeAll(async () => { - for (const key of Object.keys(testEnv)) { - previousEnv[key] = process.env[key]; + return count; } - Object.assign(process.env, testEnv); - app = buildApp().app; - const address = await app.listen({ host: "127.0.0.1", port: 0 }); - wsBase = `ws://127.0.0.1:${new URL(address).port}`; -}); + async createTextRoom(input: { + code: string; + text: string; + version: number; + pinHash?: string | null; + createdAt: Date; + updatedAt: Date; + expiresAt: Date; + }): Promise { + if (this.rooms.has(input.code)) { + return null; + } + + const row: TextRoomRow = { + code: input.code, + text: input.text, + version: input.version, + pin_hash: input.pinHash ?? null, + created_at: new Date(input.createdAt), + updated_at: new Date(input.updatedAt), + expires_at: new Date(input.expiresAt), + deleted_at: null, + }; + this.rooms.set(input.code, row); + return copyRoom(row); + } -afterAll(async () => { - await Promise.all([...openSockets].map(closeSocket)); - app.server.closeAllConnections?.(); + async findTextRoomByCode(code: string): Promise { + const room = this.rooms.get(code); + if (!room || room.deleted_at !== null) { + return null; + } - const { promise: closeTimeout, resolve: onCloseTimeout } = Promise.withResolvers(); - const timer = setTimeout(onCloseTimeout, 1000); - await Promise.race([app.close(), closeTimeout]); - clearTimeout(timer); + return copyRoom(room); + } - for (const [key, value] of Object.entries(previousEnv)) { - if (value === undefined) { - delete process.env[key]; - } else { - process.env[key] = value; + async markTextRoomActive(code: string, updatedAt: Date): Promise { + const room = this.rooms.get(code); + if (!room || room.deleted_at !== null) { + return; } + + room.updated_at = new Date(updatedAt); + room.expires_at = null; } -}); -function connect(code: string): WebSocket { - const socket = new WebSocket(`${wsBase}/api/text/${code}/ws`); - openSockets.add(socket); - socket.addEventListener("close", () => openSockets.delete(socket), { once: true }); - return socket; + async updateTextRoomText(input: { code: string; text: string; now: Date }): Promise { + const room = this.rooms.get(input.code); + if (!room || room.deleted_at !== null) { + return null; + } + + room.text = input.text; + room.version += 1; + room.updated_at = new Date(input.now); + return copyRoom(room); + } + + async scheduleTextRoomExpiry(code: string, expiresAt: Date, updatedAt: Date): Promise { + const room = this.rooms.get(code); + if (!room || room.deleted_at !== null) { + return; + } + + room.expires_at = new Date(expiresAt); + room.updated_at = new Date(updatedAt); + } + + async rearmOpenTextRooms(expiresAt: Date, updatedAt: Date): Promise { + for (const room of this.rooms.values()) { + if (room.deleted_at === null && room.expires_at === null) { + room.expires_at = new Date(expiresAt); + room.updated_at = new Date(updatedAt); + } + } + } + + async findExpiredTextRooms(now: Date, limit = 100): Promise { + const rows = [...this.rooms.values()] + .filter((room) => room.deleted_at === null && room.expires_at !== null && room.expires_at <= now) + .sort((left, right) => left.expires_at!.getTime() - right.expires_at!.getTime()) + .slice(0, limit) + .map(copyRoom); + + return rows; + } + + async markTextRoomDeleted(code: string, deletedAt: Date): Promise { + const room = this.rooms.get(code); + if (!room || room.deleted_at !== null) { + return; + } + + room.deleted_at = new Date(deletedAt); + } +} + +function copyRoom(room: TextRoomRow): TextRoomRow { + return { + code: room.code, + text: room.text, + version: room.version, + pin_hash: room.pin_hash, + created_at: new Date(room.created_at), + updated_at: new Date(room.updated_at), + expires_at: room.expires_at ? new Date(room.expires_at) : null, + deleted_at: room.deleted_at ? new Date(room.deleted_at) : null, + }; +} + +async function startTestServer(repository: TextRoomsRepository, clock: MutableClock) { + const app = Fastify({ logger: false }); + const hub = new TextSessionHub({ maxClientsPerSession: 20 }); + const sockets = new Set(); + + app.register(rateLimit, { global: false }); + app.register(websocket, { options: { maxPayload: 1024 + 1024 } }); + app.after((error) => { + if (error) { + throw error; + } + + registerTextSessionRoutes(app, { + hub, + repository, + maxBytes: 1024, + maxSessions: 500, + codeLength: 6, + ttlMs: 60 * 60 * 1000, + now: () => new Date(clock.current), + }); + }); + + const address = await app.listen({ host: "127.0.0.1", port: 0 }); + const wsBase = `ws://127.0.0.1:${new URL(address).port}`; + + return { + app, + hub, + async close() { + await Promise.all([...sockets].map(closeSocket)); + app.server.closeAllConnections?.(); + await app.close(); + }, + connect(code: string): WebSocket { + const socket = new WebSocket(`${wsBase}/api/text/${code}/ws`); + sockets.add(socket); + socket.addEventListener("close", () => sockets.delete(socket), { once: true }); + return socket; + }, + }; } function closeSocket(socket: WebSocket): Promise { @@ -91,71 +218,163 @@ function nextMessage(socket: WebSocket): Promise { return promise; } -async function createRoom(): Promise { - const response = await app.inject({ method: "POST", url: "/api/text" }); - expect(response.statusCode).toBe(200); - const payload: { code: string } = JSON.parse(response.body); - return payload.code; -} - describe("text session routes", () => { test("creates a room and serves its snapshot", async () => { - const code = await createRoom(); - expect(code).toMatch(/^[2-9A-HJ-NP-Z]{6}$/); + const clock = { current: new Date("2026-06-23T20:00:00Z") }; + const repository = new InMemoryTextRoomsRepository(); + const server = await startTestServer(repository, clock); - const snapshot = await app.inject({ method: "GET", url: `/api/text/${code}` }); - expect(snapshot.statusCode).toBe(200); - const payload: { text: string; version: number } = JSON.parse(snapshot.body); - expect(payload).toEqual({ text: "", version: 0 }); - }); + try { + const response = await server.app.inject({ method: "POST", url: "/api/text" }); + expect(response.statusCode).toBe(200); + const created: { code: string } = JSON.parse(response.body); + expect(created.code).toMatch(/^[2-9A-HJ-NP-Z]{6}$/); - test("returns 404 for an unknown room", async () => { - const response = await app.inject({ method: "GET", url: "/api/text/NOPE99" }); - expect(response.statusCode).toBe(404); - const payload: { error: string } = JSON.parse(response.body); - expect(payload.error).toBe("not_found"); + const snapshot = await server.app.inject({ method: "GET", url: `/api/text/${created.code}` }); + expect(snapshot.statusCode).toBe(200); + const payload: { text: string; version: number } = JSON.parse(snapshot.body); + expect(payload).toEqual({ text: "", version: 0 }); + } finally { + await server.close(); + } }); test("broadcasts a write to other clients and acks the writer", async () => { - const code = await createRoom(); - const author = connect(code); - const snapshot = await nextMessage(author); - expect(snapshot.type).toBe("snapshot"); - expect(snapshot.text).toBe(""); + const clock = { current: new Date("2026-06-23T20:00:00Z") }; + const repository = new InMemoryTextRoomsRepository(); + const server = await startTestServer(repository, clock); - const viewer = connect(code); - await nextMessage(viewer); + try { + const createResponse = await server.app.inject({ method: "POST", url: "/api/text" }); + const created: { code: string } = JSON.parse(createResponse.body); + const author = server.connect(created.code); + const snapshot = await nextMessage(author); + expect(snapshot.type).toBe("snapshot"); + expect(snapshot.text).toBe(""); - const updateOnViewer = nextMessage(viewer); - const ackOnAuthor = nextMessage(author); - author.send(JSON.stringify({ type: "write", text: "select 1", baseVersion: snapshot.version })); + const viewer = server.connect(created.code); + await nextMessage(viewer); - expect(await updateOnViewer).toEqual({ type: "update", text: "select 1", version: 1, by: snapshot.clientId }); - expect(await ackOnAuthor).toEqual({ type: "ack", version: 1 }); + const updateOnViewer = nextMessage(viewer); + const ackOnAuthor = nextMessage(author); + author.send(JSON.stringify({ type: "write", text: "select 1", baseVersion: snapshot.version })); - const persisted = await app.inject({ method: "GET", url: `/api/text/${code}` }); - const persistedPayload: { text: string; version: number } = JSON.parse(persisted.body); - expect(persistedPayload).toEqual({ text: "select 1", version: 1 }); + expect(await updateOnViewer).toEqual({ type: "update", text: "select 1", version: 1, by: snapshot.clientId }); + expect(await ackOnAuthor).toEqual({ type: "ack", version: 1 }); - await closeSocket(author); - await closeSocket(viewer); + const persisted = await server.app.inject({ method: "GET", url: `/api/text/${created.code}` }); + const persistedPayload: { text: string; version: number } = JSON.parse(persisted.body); + expect(persistedPayload).toEqual({ text: "select 1", version: 1 }); + } finally { + await server.close(); + } }); test("rejects an oversized write without closing the socket", async () => { - const code = await createRoom(); - const author = connect(code); - await nextMessage(author); + const clock = { current: new Date("2026-06-23T20:00:00Z") }; + const repository = new InMemoryTextRoomsRepository(); + const server = await startTestServer(repository, clock); - const errorMessage = nextMessage(author); - author.send(JSON.stringify({ type: "write", text: "x".repeat(1100), baseVersion: 0 })); + try { + const createResponse = await server.app.inject({ method: "POST", url: "/api/text" }); + const created: { code: string } = JSON.parse(createResponse.body); + const author = server.connect(created.code); + await nextMessage(author); - expect((await errorMessage).type).toBe("error"); + const errorMessage = nextMessage(author); + author.send(JSON.stringify({ type: "write", text: "x".repeat(1100), baseVersion: 0 })); - await closeSocket(author); + expect((await errorMessage).type).toBe("error"); + } finally { + await server.close(); + } }); test("rejects an unknown room code on the socket", async () => { - const socket = connect("ZZZZZZ"); - expect(await nextMessage(socket)).toEqual({ type: "error", message: "Sala não encontrada." }); + const clock = { current: new Date("2026-06-23T20:00:00Z") }; + const repository = new InMemoryTextRoomsRepository(); + const server = await startTestServer(repository, clock); + + try { + const socket = server.connect("ZZZZZZ"); + expect(await nextMessage(socket)).toEqual({ type: "error", message: "Sala não encontrada." }); + } finally { + await server.close(); + } + }); + + test("preserves room text across app restarts", async () => { + const clock = { current: new Date("2026-06-23T20:00:00Z") }; + const repository = new InMemoryTextRoomsRepository(); + const firstServer = await startTestServer(repository, clock); + + const createResponse = await firstServer.app.inject({ method: "POST", url: "/api/text" }); + const created: { code: string } = JSON.parse(createResponse.body); + const author = firstServer.connect(created.code); + await nextMessage(author); + const ack = nextMessage(author); + author.send(JSON.stringify({ type: "write", text: "survives restart", baseVersion: 0 })); + await ack; + await firstServer.close(); + + const secondServer = await startTestServer(repository, clock); + try { + const snapshot = await secondServer.app.inject({ method: "GET", url: `/api/text/${created.code}` }); + expect(snapshot.statusCode).toBe(200); + const payload: { text: string; version: number } = JSON.parse(snapshot.body); + expect(payload).toEqual({ text: "survives restart", version: 1 }); + } finally { + await secondServer.close(); + } + }); + + test("rearms previously active rooms after restart", async () => { + const clock = { current: new Date("2026-06-23T20:00:00Z") }; + const repository = new InMemoryTextRoomsRepository(); + const created = await repository.createTextRoom({ + code: "ROOM01", + text: "", + version: 0, + createdAt: clock.current, + updatedAt: clock.current, + expiresAt: new Date(clock.current.getTime() + 60 * 60 * 1000), + }); + + if (!created) { + throw new Error("expected room to be created"); + } + + await repository.markTextRoomActive(created.code, clock.current); + const active = await repository.findTextRoomByCode(created.code); + expect(active?.expires_at).toBeNull(); + + clock.current = new Date("2026-06-23T20:05:00Z"); + await rearmTextRoomsAfterRestart(60 * 60 * 1000, repository, () => new Date(clock.current)); + + const rearmed = await repository.findTextRoomByCode(created.code); + expect(rearmed?.expires_at?.toISOString()).toBe("2026-06-23T21:05:00.000Z"); + }); + + test("marks expired rooms as deleted during the database sweep", async () => { + const clock = { current: new Date("2026-06-23T20:00:00Z") }; + const repository = new InMemoryTextRoomsRepository(); + const created = await repository.createTextRoom({ + code: "ROOM01", + text: "expired", + version: 0, + createdAt: new Date("2026-06-23T18:00:00Z"), + updatedAt: new Date("2026-06-23T18:00:00Z"), + expiresAt: new Date("2026-06-23T19:00:00Z"), + }); + + if (!created) { + throw new Error("expected room to be created"); + } + + const timer = startTextSessionSweep(repository, 5, () => new Date(clock.current)); + await Bun.sleep(20); + clearInterval(timer); + + expect(await repository.findTextRoomByCode("ROOM01")).toBeNull(); }); }); diff --git a/src/server/text-session-service.ts b/src/server/text-session-service.ts index 39ee8f4..d969e1d 100644 --- a/src/server/text-session-service.ts +++ b/src/server/text-session-service.ts @@ -1,46 +1,86 @@ import { randomUUID } from "node:crypto"; import type { FastifyInstance } from "fastify"; import type { RawData, WebSocket } from "ws"; -import { normalizeSessionCode } from "./ids"; -import type { TextSessionStore } from "./text-session-store"; +import { generateSessionCode, normalizeSessionCode } from "./ids"; +import type { TextSessionHub } from "./text-session-hub"; +import type { TextRoomRow, TextRoomsRepository } from "./text-rooms-repository"; +import { textRoomsRepository } from "./text-rooms-repository"; const liveSockets = new WeakSet(); +const CODE_GENERATION_ATTEMPTS = 8; +const EXPIRED_SWEEP_LIMIT = 100; + +export type TextSessionRouteDeps = { + hub: TextSessionHub; + repository?: TextRoomsRepository; + maxBytes: number; + maxSessions: number; + codeLength: number; + ttlMs: number; + now?: () => Date; +}; + +export function registerTextSessionRoutes(app: FastifyInstance, deps: TextSessionRouteDeps): void { + const repository = deps.repository ?? textRoomsRepository; + const now = deps.now ?? (() => new Date()); -export function registerTextSessionRoutes(app: FastifyInstance, store: TextSessionStore): void { app.post( "/api/text", { preHandler: app.rateLimit({ max: 60, timeWindow: "1 hour" }) }, async (_request, reply) => { - const created = store.createSession(); - - if (!created.ok) { + if ((await repository.countActiveTextRooms(now())) >= deps.maxSessions) { reply.code(503).send({ error: "session_limit", message: "Limite de salas atingido. Tente mais tarde." }); return; } - return { code: created.session.code }; + for (let attempt = 0; attempt < CODE_GENERATION_ATTEMPTS; attempt += 1) { + const createdAt = now(); + const code = generateSessionCode(deps.codeLength); + const room = await repository.createTextRoom({ + code, + text: "", + version: 0, + createdAt, + updatedAt: createdAt, + expiresAt: new Date(createdAt.getTime() + deps.ttlMs), + }); + + if (room) { + return { code: room.code }; + } + } + + reply.code(503).send({ error: "code_exhausted", message: "Não foi possível reservar um código de sala." }); }, ); app.get<{ Params: { code: string } }>("/api/text/:code", async (request, reply) => { - const session = store.getSession(request.params.code); + const code = normalizeSessionCode(request.params.code); + const room = await repository.findTextRoomByCode(code); - if (!session) { + if (!room || isExpired(room, now(), deps.hub.clientCount(code))) { reply.code(404).send({ error: "not_found", message: "Sala não encontrada." }); return; } - return { text: session.text, version: session.version }; + return { text: room.text, version: room.version }; }); app.get<{ Params: { code: string } }>( "/api/text/:code/ws", { websocket: true }, - (socket, request) => { + async (socket, request) => { const code = normalizeSessionCode(request.params.code); - const clientId = randomUUID(); + const room = await repository.findTextRoomByCode(code); + + if (!room || isExpired(room, now(), deps.hub.clientCount(code))) { + socket.send(JSON.stringify({ type: "error", message: "Sala não encontrada." })); + socket.close(1008, "not_found"); + return; + } - const joined = store.join(code, { + const clientId = randomUUID(); + const joined = deps.hub.join(code, { id: clientId, send: (data) => { try { @@ -59,55 +99,84 @@ export function registerTextSessionRoutes(app: FastifyInstance, store: TextSessi }); if (!joined.ok) { - const message = joined.reason === "full" ? "Sala cheia." : "Sala não encontrada."; - socket.send(JSON.stringify({ type: "error", message })); + socket.send(JSON.stringify({ type: "error", message: "Sala cheia." })); socket.close(1008, joined.reason); return; } + if (joined.clientCount === 1) { + await repository.markTextRoomActive(code, now()); + } + liveSockets.add(socket); socket.on("pong", () => liveSockets.add(socket)); - - socket.send(JSON.stringify({ type: "snapshot", text: joined.text, version: joined.version, clientId })); + socket.send(JSON.stringify({ type: "snapshot", text: room.text, version: room.version, clientId })); socket.on("message", (raw: RawData) => { - const message = parseWriteMessage(raw); - - if (!message) { - return; - } - - const result = store.applyWrite(code, message.text); + void handleWrite(raw, code, clientId, socket, deps, repository, now); + }); - if (!result.ok) { - if (result.reason === "too_large") { - socket.send(JSON.stringify({ type: "error", message: "Texto excede o limite da sala." })); - } + let closed = false; + const onClose = () => { + if (closed) { return; } - const session = store.getSession(code); - if (session) { - store.broadcast( - session, - JSON.stringify({ type: "update", text: message.text, version: result.version, by: clientId }), - clientId, + closed = true; + liveSockets.delete(socket); + const remaining = deps.hub.leave(code, clientId); + + if (remaining === 0) { + const closedAt = now(); + void repository.scheduleTextRoomExpiry( + code, + new Date(closedAt.getTime() + deps.ttlMs), + closedAt, ); } - - socket.send(JSON.stringify({ type: "ack", version: result.version })); - }); - - const onClose = () => { - liveSockets.delete(socket); - store.leave(code, clientId); }; + socket.on("close", onClose); socket.on("error", onClose); }, ); } +async function handleWrite( + raw: RawData, + code: string, + clientId: string, + socket: WebSocket, + deps: TextSessionRouteDeps, + repository: TextRoomsRepository, + now: () => Date, +): Promise { + const message = parseWriteMessage(raw); + + if (!message) { + return; + } + + if (Buffer.byteLength(message.text, "utf8") > deps.maxBytes) { + socket.send(JSON.stringify({ type: "error", message: "Texto excede o limite da sala." })); + return; + } + + const updated = await repository.updateTextRoomText({ code, text: message.text, now: now() }); + + if (!updated || isExpired(updated, now(), deps.hub.clientCount(code))) { + socket.send(JSON.stringify({ type: "error", message: "Sala não encontrada." })); + return; + } + + deps.hub.broadcast( + code, + JSON.stringify({ type: "update", text: updated.text, version: updated.version, by: clientId }), + clientId, + ); + socket.send(JSON.stringify({ type: "ack", version: updated.version })); +} + function parseWriteMessage(raw: RawData): { text: string } | null { let parsed: unknown; @@ -131,19 +200,40 @@ function parseWriteMessage(raw: RawData): { text: string } | null { return null; } -export function startTextSessionSweep( - store: TextSessionStore, +function isExpired(room: TextRoomRow, now: Date, activeClientCount: number): boolean { + return activeClientCount === 0 && room.expires_at !== null && room.expires_at <= now; +} + +export async function rearmTextRoomsAfterRestart( ttlMs: number, + repository: TextRoomsRepository = textRoomsRepository, + now: () => Date = () => new Date(), +): Promise { + const reopenedAt = now(); + await repository.rearmOpenTextRooms(new Date(reopenedAt.getTime() + ttlMs), reopenedAt); +} + +export function startTextSessionSweep( + repository: TextRoomsRepository = textRoomsRepository, intervalMs = 60 * 1000, + now: () => Date = () => new Date(), ): NodeJS.Timeout { const timer = setInterval(() => { - store.sweepExpired(ttlMs); + void sweepExpiredRooms(repository, now()); }, intervalMs); timer.unref(); return timer; } +async function sweepExpiredRooms(repository: TextRoomsRepository, currentTime: Date): Promise { + const expiredRooms = await repository.findExpiredTextRooms(currentTime, EXPIRED_SWEEP_LIMIT); + + for (const room of expiredRooms) { + await repository.markTextRoomDeleted(room.code, currentTime); + } +} + export function startTextSessionHeartbeat(app: FastifyInstance, intervalMs = 30 * 1000): NodeJS.Timeout { const timer = setInterval(() => { for (const socket of app.websocketServer.clients) { diff --git a/src/server/text-session-store.test.ts b/src/server/text-session-store.test.ts index b922308..8c9e1c6 100644 --- a/src/server/text-session-store.test.ts +++ b/src/server/text-session-store.test.ts @@ -21,125 +21,48 @@ function makeClient(id: string): RecordingClient { function makeStore(overrides: Partial = {}) { return new TextSessionStore({ - maxBytes: 1024, - maxSessions: 100, maxClientsPerSession: 10, - codeLength: 6, - generateCode: () => "AAAAAA", ...overrides, }); } describe("TextSessionStore", () => { - test("creates an empty session and tracks size", () => { + test("joins clients into a room and tracks the count", () => { const store = makeStore(); - const result = store.createSession(); - expect(result.ok).toBe(true); - if (result.ok) { - expect(result.session.text).toBe(""); - expect(result.session.version).toBe(0); - expect(result.session.code).toBe("AAAAAA"); - } - expect(store.size).toBe(1); + expect(store.join("ROOM01", makeClient("a"))).toEqual({ ok: true, clientCount: 1 }); + expect(store.join("ROOM01", makeClient("b"))).toEqual({ ok: true, clientCount: 2 }); + expect(store.clientCount("room01")).toBe(2); }); - test("retries code generation on collision", () => { - const codes = ["AAAAAA", "AAAAAA", "BBBBBB"]; - let index = 0; - const store = makeStore({ generateCode: () => codes[index++]! }); + test("enforces the per-room client cap", () => { + const store = makeStore({ maxClientsPerSession: 1 }); - const first = store.createSession(); - const second = store.createSession(); - - expect(first.ok && first.session.code).toBe("AAAAAA"); - expect(second.ok && second.session.code).toBe("BBBBBB"); - expect(store.size).toBe(2); - }); - - test("rejects creation past the session limit", () => { - const codes = ["AAAAAA", "BBBBBB"]; - let index = 0; - const store = makeStore({ maxSessions: 1, generateCode: () => codes[index++]! }); - - expect(store.createSession().ok).toBe(true); - const second = store.createSession(); - - expect(second).toEqual({ ok: false, reason: "limit" }); - expect(store.size).toBe(1); - }); - - test("looks up sessions case-insensitively", () => { - const store = makeStore({ generateCode: () => "ABCDEF" }); - store.createSession(); - - expect(store.getSession("abcdef")?.code).toBe("ABCDEF"); - expect(store.getSession(" abcDef ")?.code).toBe("ABCDEF"); - }); - - test("join returns a snapshot and enforces the client cap", () => { - const store = makeStore({ maxClientsPerSession: 1, generateCode: () => "ROOM01" }); - store.createSession(); - - expect(store.join("ROOM01", makeClient("a"))).toEqual({ ok: true, text: "", version: 0 }); + expect(store.join("ROOM01", makeClient("a"))).toEqual({ ok: true, clientCount: 1 }); expect(store.join("ROOM01", makeClient("b"))).toEqual({ ok: false, reason: "full" }); - expect(store.join("NOPE00", makeClient("c"))).toEqual({ ok: false, reason: "not_found" }); - }); - - test("applyWrite bumps the version and stores the text", () => { - const store = makeStore({ generateCode: () => "ROOM01" }); - store.createSession(); - - expect(store.applyWrite("ROOM01", "select 1")).toEqual({ ok: true, version: 1 }); - expect(store.applyWrite("room01", "select 2")).toEqual({ ok: true, version: 2 }); - expect(store.getSession("ROOM01")?.text).toBe("select 2"); - expect(store.applyWrite("NOPE00", "x")).toEqual({ ok: false, reason: "not_found" }); - }); - - test("applyWrite rejects oversized text without mutating the session", () => { - const store = makeStore({ maxBytes: 5, generateCode: () => "ROOM01" }); - store.createSession(); - store.applyWrite("ROOM01", "okay"); - - expect(store.applyWrite("ROOM01", "way too long")).toEqual({ ok: false, reason: "too_large" }); - expect(store.getSession("ROOM01")?.text).toBe("okay"); - expect(store.getSession("ROOM01")?.version).toBe(1); }); test("broadcast delivers to every client except the excluded one", () => { - const store = makeStore({ generateCode: () => "ROOM01" }); - store.createSession(); + const store = makeStore(); const author = makeClient("author"); const viewer = makeClient("viewer"); store.join("ROOM01", author); store.join("ROOM01", viewer); - const session = store.getSession("ROOM01")!; - store.broadcast(session, "update", "author"); + store.broadcast("ROOM01", "update", "author"); expect(author.sent).toEqual([]); expect(viewer.sent).toEqual(["update"]); }); - test("sweepExpired removes only idle, empty sessions past the TTL", () => { - let clock = 0; - const codes = ["IDLE00", "BUSY00", "FRESH0"]; - let index = 0; - const store = makeStore({ generateCode: () => codes[index++]!, now: () => clock }); - - store.createSession(); // IDLE00 at t=0, no clients - store.createSession(); // BUSY00 at t=0 - store.join("BUSY00", makeClient("a")); // keeps BUSY00 alive - - clock = 5000; - store.createSession(); // FRESH0 at t=5000 - - clock = 6000; - const removed = store.sweepExpired(2000); + test("leave removes a client and deletes the empty room", () => { + const store = makeStore(); + store.join("ROOM01", makeClient("a")); + store.join("ROOM01", makeClient("b")); - expect(removed).toBe(1); - expect(store.getSession("IDLE00")).toBeUndefined(); - expect(store.getSession("BUSY00")).toBeDefined(); - expect(store.getSession("FRESH0")).toBeDefined(); + expect(store.leave("ROOM01", "a")).toBe(1); + expect(store.clientCount("ROOM01")).toBe(1); + expect(store.leave("ROOM01", "b")).toBe(0); + expect(store.clientCount("ROOM01")).toBe(0); }); }); diff --git a/src/server/text-session-store.ts b/src/server/text-session-store.ts index 51d57bc..9c1c620 100644 --- a/src/server/text-session-store.ts +++ b/src/server/text-session-store.ts @@ -1,163 +1,6 @@ -import { generateSessionCode, normalizeSessionCode } from "./ids"; - -export type SessionClient = { - id: string; - send: (data: string) => void; - close: (code?: number, reason?: string) => void; -}; - -export type TextSession = { - code: string; - text: string; - version: number; - updatedAt: number; - clients: Map; -}; - -export type TextSessionStoreOptions = { - maxBytes: number; - maxSessions: number; - maxClientsPerSession: number; - codeLength: number; - generateCode?: (length: number) => string; - now?: () => number; -}; - -export type CreateResult = - | { ok: true; session: TextSession } - | { ok: false; reason: "limit" }; - -export type JoinResult = - | { ok: true; text: string; version: number } - | { ok: false; reason: "not_found" | "full" }; - -export type WriteResult = - | { ok: true; version: number } - | { ok: false; reason: "not_found" | "too_large" }; - -const CODE_GENERATION_ATTEMPTS = 8; - -export class TextSessionStore { - private readonly sessions = new Map(); - private readonly maxBytes: number; - private readonly maxSessions: number; - private readonly maxClientsPerSession: number; - private readonly codeLength: number; - private readonly generateCode: (length: number) => string; - private readonly now: () => number; - - constructor(options: TextSessionStoreOptions) { - this.maxBytes = options.maxBytes; - this.maxSessions = options.maxSessions; - this.maxClientsPerSession = options.maxClientsPerSession; - this.codeLength = options.codeLength; - this.generateCode = options.generateCode ?? generateSessionCode; - this.now = options.now ?? Date.now; - } - - get size(): number { - return this.sessions.size; - } - - createSession(): CreateResult { - if (this.sessions.size >= this.maxSessions) { - return { ok: false, reason: "limit" }; - } - - for (let attempt = 0; attempt < CODE_GENERATION_ATTEMPTS; attempt += 1) { - const code = normalizeSessionCode(this.generateCode(this.codeLength)); - - if (!this.sessions.has(code)) { - const session: TextSession = { - code, - text: "", - version: 0, - updatedAt: this.now(), - clients: new Map(), - }; - this.sessions.set(code, session); - return { ok: true, session }; - } - } - - return { ok: false, reason: "limit" }; - } - - getSession(code: string): TextSession | undefined { - return this.sessions.get(normalizeSessionCode(code)); - } - - join(code: string, client: SessionClient): JoinResult { - const session = this.getSession(code); - - if (!session) { - return { ok: false, reason: "not_found" }; - } - - if (session.clients.size >= this.maxClientsPerSession) { - return { ok: false, reason: "full" }; - } - - session.clients.set(client.id, client); - session.updatedAt = this.now(); - - return { ok: true, text: session.text, version: session.version }; - } - - leave(code: string, clientId: string): void { - const session = this.getSession(code); - - if (!session) { - return; - } - - session.clients.delete(clientId); - session.updatedAt = this.now(); - } - - applyWrite(code: string, text: string): WriteResult { - const session = this.getSession(code); - - if (!session) { - return { ok: false, reason: "not_found" }; - } - - if (Buffer.byteLength(text, "utf8") > this.maxBytes) { - return { ok: false, reason: "too_large" }; - } - - session.text = text; - session.version += 1; - session.updatedAt = this.now(); - - return { ok: true, version: session.version }; - } - - broadcast(session: TextSession, data: string, exceptId?: string): void { - for (const client of session.clients.values()) { - if (client.id === exceptId) { - continue; - } - - try { - client.send(data); - } catch { - // A failing client must not break delivery to the rest of the room. - } - } - } - - sweepExpired(ttlMs: number): number { - const now = this.now(); - let removed = 0; - - for (const [code, session] of this.sessions) { - if (session.clients.size === 0 && now - session.updatedAt > ttlMs) { - this.sessions.delete(code); - removed += 1; - } - } - - return removed; - } -} +export { + TextSessionHub as TextSessionStore, + type JoinResult, + type SessionClient, + type TextSessionHubOptions as TextSessionStoreOptions, +} from "./text-session-hub";