Complete Planner session management (tasks 4.2, 4.6, 4.7) (#9)
This commit is contained in:
parent
a9a6b27af4
commit
6a3e566438
9 changed files with 201 additions and 40 deletions
23
apps/planner/app/lib/db.ts
Normal file
23
apps/planner/app/lib/db.ts
Normal file
|
|
@ -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
|
||||
)
|
||||
`;
|
||||
}
|
||||
|
|
@ -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<string, SessionMetadata>();
|
||||
|
||||
export function createSession(options?: {
|
||||
export async function createSession(options?: {
|
||||
callbackUrl?: string;
|
||||
callbackToken?: string;
|
||||
}): SessionMetadata {
|
||||
}): Promise<SessionMetadata> {
|
||||
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<SessionMetadata | undefined> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
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<Uint8Array | null> {
|
||||
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<boolean> {
|
||||
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<SessionMetadata[]> {
|
||||
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<number> {
|
||||
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,
|
||||
};
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<string, Y.Doc>();
|
||||
const sessionClients = new Map<string, Set<WebSocket>>();
|
||||
|
||||
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<Y.Doc> {
|
||||
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<string, ReturnType<typeof setTimeout>>();
|
||||
|
||||
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(() => {});
|
||||
}
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
|
|
|
|||
|
|
@ -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 });
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 });
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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:",
|
||||
|
|
|
|||
19
docker-compose.dev.yml
Normal file
19
docker-compose.dev.yml
Normal file
|
|
@ -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:
|
||||
|
|
@ -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
|
||||
|
|
|
|||
9
pnpm-lock.yaml
generated
9
pnpm-lock.yaml
generated
|
|
@ -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: {}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue