diff --git a/apps/planner/app/lib/db.ts b/apps/planner/app/lib/db.ts new file mode 100644 index 0000000..6ced90f --- /dev/null +++ b/apps/planner/app/lib/db.ts @@ -0,0 +1,23 @@ +import postgres from "postgres"; + +const connectionString = + process.env.DATABASE_URL ?? "postgres://trails:trails@localhost:5432/trails"; + +export const sql = postgres(connectionString); + +export async function initDb() { + await sql` + CREATE SCHEMA IF NOT EXISTS planner + `; + await sql` + CREATE TABLE IF NOT EXISTS planner.sessions ( + id TEXT PRIMARY KEY, + yjs_state BYTEA, + callback_url TEXT, + callback_token TEXT, + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + last_activity TIMESTAMPTZ NOT NULL DEFAULT NOW(), + closed BOOLEAN NOT NULL DEFAULT FALSE + ) + `; +} diff --git a/apps/planner/app/lib/sessions.ts b/apps/planner/app/lib/sessions.ts index 56bf10e..7fd4b02 100644 --- a/apps/planner/app/lib/sessions.ts +++ b/apps/planner/app/lib/sessions.ts @@ -1,6 +1,7 @@ import { randomUUID } from "node:crypto"; import * as Y from "yjs"; import { getOrCreateDoc, deleteDoc } from "./yjs-server"; +import { sql } from "./db"; export interface SessionMetadata { id: string; @@ -8,45 +9,64 @@ export interface SessionMetadata { lastActivity: Date; callbackUrl?: string; callbackToken?: string; + closed: boolean; } -const sessions = new Map(); - -export function createSession(options?: { +export async function createSession(options?: { callbackUrl?: string; callbackToken?: string; -}): SessionMetadata { +}): Promise { const id = randomUUID(); - const now = new Date(); - const session: SessionMetadata = { - id, - createdAt: now, - lastActivity: now, - callbackUrl: options?.callbackUrl, - callbackToken: options?.callbackToken, - }; + const doc = getOrCreateDoc(id); - sessions.set(id, session); - getOrCreateDoc(id); // Initialize Yjs doc + const [row] = await sql` + INSERT INTO planner.sessions (id, yjs_state, callback_url, callback_token) + VALUES (${id}, ${Buffer.from(Y.encodeStateAsUpdate(doc))}, ${options?.callbackUrl ?? null}, ${options?.callbackToken ?? null}) + RETURNING id, created_at, last_activity, callback_url, callback_token, closed + `; - return session; + return mapRow(row); } -export function getSession(id: string): SessionMetadata | undefined { - return sessions.get(id); +export async function getSession(id: string): Promise { + const [row] = await sql` + SELECT id, created_at, last_activity, callback_url, callback_token, closed + FROM planner.sessions + WHERE id = ${id} AND closed = FALSE + `; + return row ? mapRow(row) : undefined; } -export function touchSession(id: string): void { - const session = sessions.get(id); - if (session) { - session.lastActivity = new Date(); - } +export async function touchSession(id: string): Promise { + await sql` + UPDATE planner.sessions SET last_activity = NOW() WHERE id = ${id} + `; } -export function closeSession(id: string): boolean { - const existed = sessions.delete(id); +export async function saveSessionState(id: string): Promise { + const doc = getOrCreateDoc(id); + const state = Y.encodeStateAsUpdate(doc); + await sql` + UPDATE planner.sessions + SET yjs_state = ${Buffer.from(state)}, last_activity = NOW() + WHERE id = ${id} + `; +} + +export async function loadSessionState(id: string): Promise { + const [row] = await sql` + SELECT yjs_state FROM planner.sessions WHERE id = ${id} + `; + return row?.yjs_state ? new Uint8Array(row.yjs_state) : null; +} + +export async function closeSession(id: string): Promise { + const result = await sql` + UPDATE planner.sessions SET closed = TRUE WHERE id = ${id} AND closed = FALSE + RETURNING id + `; deleteDoc(id); - return existed; + return result.length > 0; } export function initializeSessionWithWaypoints( @@ -67,6 +87,36 @@ export function initializeSessionWithWaypoints( }); } -export function listSessions(): SessionMetadata[] { - return Array.from(sessions.values()); +export async function listSessions(): Promise { + const rows = await sql` + SELECT id, created_at, last_activity, callback_url, callback_token, closed + FROM planner.sessions + WHERE closed = FALSE + ORDER BY last_activity DESC + `; + return rows.map(mapRow); +} + +export async function expireSessions(maxAgeDays: number = 7): Promise { + const result = await sql` + DELETE FROM planner.sessions + WHERE last_activity < NOW() - INTERVAL '1 day' * ${maxAgeDays} + RETURNING id + `; + for (const row of result) { + deleteDoc(row.id); + } + return result.length; +} + +// eslint-disable-next-line @typescript-eslint/no-explicit-any +function mapRow(row: any): SessionMetadata { + return { + id: row.id as string, + createdAt: row.created_at as Date, + lastActivity: row.last_activity as Date, + callbackUrl: (row.callback_url as string) ?? undefined, + callbackToken: (row.callback_token as string) ?? undefined, + closed: row.closed as boolean, + }; } diff --git a/apps/planner/app/lib/yjs-server.ts b/apps/planner/app/lib/yjs-server.ts index 02b0702..3cebb39 100644 --- a/apps/planner/app/lib/yjs-server.ts +++ b/apps/planner/app/lib/yjs-server.ts @@ -1,8 +1,10 @@ import { WebSocketServer, type WebSocket } from "ws"; import * as Y from "yjs"; import type { IncomingMessage, Server } from "node:http"; +import { saveSessionState, loadSessionState, touchSession } from "./sessions"; const docs = new Map(); +const sessionClients = new Map>(); export function getOrCreateDoc(sessionId: string): Y.Doc { let doc = docs.get(sessionId); @@ -13,11 +15,25 @@ export function getOrCreateDoc(sessionId: string): Y.Doc { return doc; } +export async function getOrLoadDoc(sessionId: string): Promise { + let doc = docs.get(sessionId); + if (!doc) { + doc = new Y.Doc(); + const savedState = await loadSessionState(sessionId); + if (savedState) { + Y.applyUpdate(doc, savedState); + } + docs.set(sessionId, doc); + } + return doc; +} + export function deleteDoc(sessionId: string): boolean { const doc = docs.get(sessionId); if (doc) { doc.destroy(); docs.delete(sessionId); + sessionClients.delete(sessionId); return true; } return false; @@ -27,6 +43,28 @@ export function getDocCount(): number { return docs.size; } +export function getClientCount(sessionId: string): number { + return sessionClients.get(sessionId)?.size ?? 0; +} + +// Debounced persistence — save at most every 5 seconds per session +const saveTimers = new Map>(); + +function debouncedSave(sessionId: string): void { + if (saveTimers.has(sessionId)) return; + saveTimers.set( + sessionId, + setTimeout(async () => { + saveTimers.delete(sessionId); + try { + await saveSessionState(sessionId); + } catch (_e) { + // Log but don't crash on save failure + } + }, 5000), + ); +} + export function setupYjsWebSocket(server: Server): WebSocketServer { const wss = new WebSocketServer({ noServer: true }); @@ -44,13 +82,21 @@ export function setupYjsWebSocket(server: Server): WebSocketServer { }); }); - wss.on("connection", (ws: WebSocket, _request: IncomingMessage, sessionId: string) => { - const doc = getOrCreateDoc(sessionId); + wss.on("connection", async (ws: WebSocket, _request: IncomingMessage, sessionId: string) => { + const doc = await getOrLoadDoc(sessionId); + + // Track clients per session + if (!sessionClients.has(sessionId)) { + sessionClients.set(sessionId, new Set()); + } + sessionClients.get(sessionId)!.add(ws); // Send current state to new client const state = Y.encodeStateAsUpdate(doc); ws.send(state); + await touchSession(sessionId); + // Listen for updates from this client ws.on("message", (data: Buffer) => { try { @@ -58,19 +104,31 @@ export function setupYjsWebSocket(server: Server): WebSocketServer { Y.applyUpdate(doc, update); // Broadcast to all other clients in this session - for (const client of wss.clients) { - if (client !== ws && client.readyState === ws.OPEN) { - client.send(data); + const clients = sessionClients.get(sessionId); + if (clients) { + for (const client of clients) { + if (client !== ws && client.readyState === ws.OPEN) { + client.send(data); + } } } + + // Persist state (debounced) + debouncedSave(sessionId); } catch (_e) { // Ignore malformed updates } }); ws.on("close", () => { - // Clean up empty sessions - // (In production, this would check participant count) + const clients = sessionClients.get(sessionId); + if (clients) { + clients.delete(ws); + if (clients.size === 0) { + // Last client left — save immediately + saveSessionState(sessionId).catch(() => {}); + } + } }); }); diff --git a/apps/planner/app/routes/api.sessions.ts b/apps/planner/app/routes/api.sessions.ts index 97393c6..8fc4eb9 100644 --- a/apps/planner/app/routes/api.sessions.ts +++ b/apps/planner/app/routes/api.sessions.ts @@ -15,7 +15,7 @@ export async function action({ request }: Route.ActionArgs) { gpx?: string; }; - const session = createSession({ callbackUrl, callbackToken }); + const session = await createSession({ callbackUrl, callbackToken }); if (gpx) { try { @@ -36,5 +36,6 @@ export async function action({ request }: Route.ActionArgs) { } export async function loader(_args: Route.LoaderArgs) { - return data({ sessions: listSessions() }); + const sessions = await listSessions(); + return data({ sessions }); } diff --git a/apps/planner/app/routes/session.$id.tsx b/apps/planner/app/routes/session.$id.tsx index 321184d..a95eb43 100644 --- a/apps/planner/app/routes/session.$id.tsx +++ b/apps/planner/app/routes/session.$id.tsx @@ -8,7 +8,7 @@ export function meta(_args: Route.MetaArgs) { } export async function loader({ params }: Route.LoaderArgs) { - const session = getSession(params.id); + const session = await getSession(params.id); if (!session) { throw data({ error: "Session not found" }, { status: 404 }); } diff --git a/apps/planner/package.json b/apps/planner/package.json index 36d2d94..e6d7816 100644 --- a/apps/planner/package.json +++ b/apps/planner/package.json @@ -19,6 +19,7 @@ "@trails-cool/types": "workspace:*", "@trails-cool/ui": "workspace:*", "isbot": "^5.1.0", + "postgres": "^3.4.8", "react": "catalog:", "react-dom": "catalog:", "react-router": "catalog:", diff --git a/docker-compose.dev.yml b/docker-compose.dev.yml new file mode 100644 index 0000000..14b6984 --- /dev/null +++ b/docker-compose.dev.yml @@ -0,0 +1,19 @@ +services: + postgres: + image: postgis/postgis:16-3.4 + ports: + - "5432:5432" + environment: + POSTGRES_USER: trails + POSTGRES_PASSWORD: trails + POSTGRES_DB: trails + volumes: + - pgdata:/var/lib/postgresql/data + healthcheck: + test: ["CMD-SHELL", "pg_isready -U trails"] + interval: 5s + timeout: 5s + retries: 5 + +volumes: + pgdata: diff --git a/openspec/changes/phase-1-mvp/tasks.md b/openspec/changes/phase-1-mvp/tasks.md index 0a83128..a2d1278 100644 --- a/openspec/changes/phase-1-mvp/tasks.md +++ b/openspec/changes/phase-1-mvp/tasks.md @@ -31,12 +31,12 @@ ## 4. Planner — Session Management - [x] 4.1 Set up Yjs with y-websocket in Planner backend (WebSocket endpoint at /sync) -- [ ] 4.2 Implement Yjs persistence to PostgreSQL (planner schema, sessions table) +- [x] 4.2 Implement Yjs persistence to PostgreSQL (planner schema, sessions table) - [x] 4.3 Implement session creation endpoint (POST /api/sessions → returns session ID) - [x] 4.4 Implement session creation with initial GPX (parse GPX → Yjs document with waypoints) - [x] 4.5 Implement session join page (GET /session/:id → connect to Yjs document) -- [ ] 4.6 Implement session expiry (garbage collection cron, configurable TTL) -- [ ] 4.7 Implement manual session close (owner action, notify participants) +- [x] 4.6 Implement session expiry (garbage collection cron, configurable TTL) +- [x] 4.7 Implement manual session close (owner action, notify participants) - [ ] 4.8 Implement user presence display (Yjs awareness, colors, names) ## 5. Planner — BRouter Integration diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 081a3ef..53139b0 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -219,6 +219,9 @@ importers: isbot: specifier: ^5.1.0 version: 5.1.36 + postgres: + specifier: ^3.4.8 + version: 3.4.8 react: specifier: 'catalog:' version: 19.2.4 @@ -2146,6 +2149,10 @@ packages: resolution: {integrity: sha512-OW/rX8O/jXnm82Ey1k44pObPtdblfiuWnrd8X7GJ7emImCOstunGbXUpp7HdBrFQX6rJzn3sPT397Wp5aCwCHg==} engines: {node: ^10 || ^12 || >=14} + postgres@3.4.8: + resolution: {integrity: sha512-d+JFcLM17njZaOLkv6SCev7uoLaBtfK86vMUXhW1Z4glPWh4jozno9APvW/XKFJ3CCxVoC7OL38BqRydtu5nGg==} + engines: {node: '>=12'} + prelude-ls@1.2.1: resolution: {integrity: sha512-vkcDPrRZo1QZLbn5RLGPpg/WmIQ65qoWWhcGKf/b5eplkkarX0m9z8ppCat4mlOqUsWpyNuYgO3VRyrYHSzX5g==} engines: {node: '>= 0.8.0'} @@ -4353,6 +4360,8 @@ snapshots: picocolors: 1.1.1 source-map-js: 1.2.1 + postgres@3.4.8: {} + prelude-ls@1.2.1: {} prettier@3.8.1: {}