WebSocket receiver
This commit is contained in:
Generated
+2
-1
@@ -61,7 +61,8 @@
|
|||||||
"react-dom": "19.2.3",
|
"react-dom": "19.2.3",
|
||||||
"recharts": "^2.15.4",
|
"recharts": "^2.15.4",
|
||||||
"tailwind-merge": "^3.4.0",
|
"tailwind-merge": "^3.4.0",
|
||||||
"tailwindcss-animate": "^1.0.7"
|
"tailwindcss-animate": "^1.0.7",
|
||||||
|
"ws": "^8.19.0"
|
||||||
},
|
},
|
||||||
"devDependencies": {
|
"devDependencies": {
|
||||||
"@tailwindcss/postcss": "^4",
|
"@tailwindcss/postcss": "^4",
|
||||||
|
|||||||
+4
-3
@@ -3,9 +3,9 @@
|
|||||||
"version": "0.1.0",
|
"version": "0.1.0",
|
||||||
"private": true,
|
"private": true,
|
||||||
"scripts": {
|
"scripts": {
|
||||||
"dev": "next dev",
|
"dev": "node server.mjs",
|
||||||
"build": "next build",
|
"build": "next build",
|
||||||
"start": "next start"
|
"start": "NODE_ENV=production node server.mjs"
|
||||||
},
|
},
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"@codemirror/lang-go": "^6.0.1",
|
"@codemirror/lang-go": "^6.0.1",
|
||||||
@@ -61,7 +61,8 @@
|
|||||||
"react-dom": "19.2.3",
|
"react-dom": "19.2.3",
|
||||||
"recharts": "^2.15.4",
|
"recharts": "^2.15.4",
|
||||||
"tailwind-merge": "^3.4.0",
|
"tailwind-merge": "^3.4.0",
|
||||||
"tailwindcss-animate": "^1.0.7"
|
"tailwindcss-animate": "^1.0.7",
|
||||||
|
"ws": "^8.19.0"
|
||||||
},
|
},
|
||||||
"devDependencies": {
|
"devDependencies": {
|
||||||
"@tailwindcss/postcss": "^4",
|
"@tailwindcss/postcss": "^4",
|
||||||
|
|||||||
+60
@@ -0,0 +1,60 @@
|
|||||||
|
import { createServer } from 'http';
|
||||||
|
import next from 'next';
|
||||||
|
import { WebSocketServer } from 'ws';
|
||||||
|
|
||||||
|
const dev = process.env.NODE_ENV !== 'production';
|
||||||
|
const app = next({ dev });
|
||||||
|
const handle = app.getRequestHandler();
|
||||||
|
|
||||||
|
app.prepare().then(() => {
|
||||||
|
const server = createServer((req, res) => {
|
||||||
|
handle(req, res);
|
||||||
|
});
|
||||||
|
|
||||||
|
const clients = {};
|
||||||
|
|
||||||
|
const wss = new WebSocketServer({ noServer: true });
|
||||||
|
|
||||||
|
// Handle WebSocket connections
|
||||||
|
wss.on('connection', (ws, req) => {
|
||||||
|
ws.on('message', (msg) => {
|
||||||
|
const runId = msg.toString();
|
||||||
|
if (!clients[runId]) clients[runId] = []
|
||||||
|
clients[runId].push(ws);
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
// Handle POST requests
|
||||||
|
server.on('request', async (req, res) => {
|
||||||
|
console.log(req.method, req.url)
|
||||||
|
if (req.method === 'POST' && req.url === '/api/run/move') {
|
||||||
|
let body = '';
|
||||||
|
req.on('data', chunk => body += chunk);
|
||||||
|
req.on('end', () => {
|
||||||
|
const move = JSON.parse(body);
|
||||||
|
|
||||||
|
clients[move.runId]?.forEach(client => {
|
||||||
|
if (client.readyState === client.OPEN) {
|
||||||
|
client.send(JSON.stringify(move.moves));
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
res.writeHead(200, { 'Content-Type': 'application/json' });
|
||||||
|
res.end(JSON.stringify({ status: 'ok' }));
|
||||||
|
});
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
// Upgrade HTTP connections to WebSocket when needed
|
||||||
|
server.on('upgrade', (req, socket, head) => {
|
||||||
|
if (req.url === '/api/view-run') {
|
||||||
|
wss.handleUpgrade(req, socket, head, (ws) => {
|
||||||
|
wss.emit('connection', ws, req);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
server.listen(3000, () => {
|
||||||
|
console.log('Server running on http://localhost:3000');
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -8,7 +8,7 @@ import { go } from "@codemirror/lang-go";
|
|||||||
import { json } from "@codemirror/lang-json";
|
import { json } from "@codemirror/lang-json";
|
||||||
import { python } from "@codemirror/lang-python";
|
import { python } from "@codemirror/lang-python";
|
||||||
import { java } from "@codemirror/lang-java";
|
import { java } from "@codemirror/lang-java";
|
||||||
import { useEffect, useState } from "react";
|
import { useEffect, useRef, useState } from "react";
|
||||||
import useAction from "@/hook/use-action";
|
import useAction from "@/hook/use-action";
|
||||||
import { getLiveRunMoves } from "@/modules/live-run/actions";
|
import { getLiveRunMoves } from "@/modules/live-run/actions";
|
||||||
import { LiveRun, LiveRunMove } from "@/modules/live-run/types";
|
import { LiveRun, LiveRunMove } from "@/modules/live-run/types";
|
||||||
@@ -43,66 +43,70 @@ function getLanguageExtension(lang: Language) {
|
|||||||
|
|
||||||
export default function Editor({ run }: { run: LiveRun | null | undefined }) {
|
export default function Editor({ run }: { run: LiveRun | null | undefined }) {
|
||||||
const [editorView, setEditorView] = useState<EditorView>();
|
const [editorView, setEditorView] = useState<EditorView>();
|
||||||
const [rawMoves] = useAction<LiveRunMove[]>(
|
const [isConnected, setConnected] = useState(false);
|
||||||
async () => (run ? await getLiveRunMoves(run.id) : []),
|
|
||||||
[run],
|
|
||||||
);
|
|
||||||
|
|
||||||
const [_, setLastMoveId] = useState<number>(-1);
|
const [rawMoves, setRawMoves] = useState<LiveRunMove[]>([]);
|
||||||
const [moves, setMoves] = useState<LiveRunMove[]>([]);
|
|
||||||
const [movesLoading, setMovesLoading] = useState(false);
|
const [movesLoading, setMovesLoading] = useState(false);
|
||||||
|
const scheduledTimeouts = useRef<NodeJS.Timeout[]>([]);
|
||||||
|
|
||||||
|
// WebSocket connection
|
||||||
useEffect(() => {
|
useEffect(() => {
|
||||||
if (!rawMoves || rawMoves.length === 0 || movesLoading) return;
|
if (!run) return;
|
||||||
|
|
||||||
setLastMoveId((prevLastMoveId) => {
|
const ws = new WebSocket(`ws://${window.location.host}/api/view-run`);
|
||||||
const newMoves = rawMoves.filter((move) => move.moveId > prevLastMoveId);
|
|
||||||
|
|
||||||
if (newMoves.length === 0) return prevLastMoveId;
|
ws.onopen = () => {
|
||||||
|
ws.send(run.id);
|
||||||
|
};
|
||||||
|
|
||||||
setMoves(newMoves);
|
ws.onmessage = (m) => {
|
||||||
|
const data: LiveRunMove[] = JSON.parse(m.data);
|
||||||
|
setRawMoves(data);
|
||||||
|
};
|
||||||
|
|
||||||
return Math.max(...newMoves.map((m) => m.moveId));
|
return () => ws.close();
|
||||||
});
|
}, [run]);
|
||||||
}, [rawMoves, movesLoading]);
|
|
||||||
|
|
||||||
|
// Scheduler
|
||||||
useEffect(() => {
|
useEffect(() => {
|
||||||
if (!editorView || !moves) return;
|
if (!editorView || !rawMoves.length) return;
|
||||||
|
|
||||||
let cancelled = false;
|
setMovesLoading(true);
|
||||||
|
|
||||||
const executeMove = async (move: LiveRunMove) => {
|
// Clear any existing scheduled timeouts
|
||||||
return new Promise<void>((resolve) => {
|
scheduledTimeouts.current.forEach((t) => clearTimeout(t));
|
||||||
|
scheduledTimeouts.current = [];
|
||||||
|
|
||||||
|
const startTime = Date.now();
|
||||||
|
|
||||||
|
rawMoves.forEach((move) => {
|
||||||
const timeout = setTimeout(() => {
|
const timeout = setTimeout(() => {
|
||||||
if (!editorView || cancelled) return resolve();
|
if (!editorView) return;
|
||||||
|
|
||||||
editorView.dispatch({
|
editorView.dispatch({
|
||||||
selection: { anchor: move.cursor },
|
selection: { anchor: move.cursor },
|
||||||
scrollIntoView: true,
|
scrollIntoView: true,
|
||||||
changes: move.changes,
|
changes: move.changes,
|
||||||
});
|
});
|
||||||
|
|
||||||
resolve();
|
|
||||||
}, move.latency);
|
}, move.latency);
|
||||||
|
|
||||||
return () => clearTimeout(timeout);
|
scheduledTimeouts.current.push(timeout);
|
||||||
});
|
});
|
||||||
};
|
|
||||||
|
|
||||||
(async () => {
|
// Stop loading after the last move
|
||||||
if (!moves) return;
|
const lastMoveLatency = rawMoves[rawMoves.length - 1].latency;
|
||||||
setMovesLoading(true);
|
const finishTimeout = setTimeout(
|
||||||
for (const move of moves) {
|
() => setMovesLoading(false),
|
||||||
if (cancelled) break;
|
lastMoveLatency + 50,
|
||||||
await executeMove(move);
|
);
|
||||||
}
|
scheduledTimeouts.current.push(finishTimeout);
|
||||||
setMovesLoading(false);
|
|
||||||
})();
|
|
||||||
|
|
||||||
return () => {
|
return () => {
|
||||||
cancelled = true;
|
scheduledTimeouts.current.forEach((t) => clearTimeout(t));
|
||||||
|
scheduledTimeouts.current = [];
|
||||||
|
setMovesLoading(false);
|
||||||
};
|
};
|
||||||
}, [moves]);
|
}, [rawMoves, editorView]);
|
||||||
|
|
||||||
return (
|
return (
|
||||||
<CodeMirror
|
<CodeMirror
|
||||||
|
|||||||
@@ -16,7 +16,7 @@ import { useParams } from "next/navigation";
|
|||||||
export default function Run() {
|
export default function Run() {
|
||||||
const params = useParams<{ runId: string }>();
|
const params = useParams<{ runId: string }>();
|
||||||
|
|
||||||
const [run] = useActionInterval(() => getLiveRun(params.runId), 1000, [
|
const [run] = useActionInterval(() => getLiveRun(params.runId), 5000, [
|
||||||
params,
|
params,
|
||||||
]);
|
]);
|
||||||
|
|
||||||
|
|||||||
@@ -4,7 +4,6 @@ import "server-only";
|
|||||||
import { LiveRun, LiveRunMove } from "./types";
|
import { LiveRun, LiveRunMove } from "./types";
|
||||||
import { SupabaseClient } from "@supabase/supabase-js";
|
import { SupabaseClient } from "@supabase/supabase-js";
|
||||||
import { redis } from "@/lib/redis";
|
import { redis } from "@/lib/redis";
|
||||||
import { randomInt } from "crypto";
|
|
||||||
|
|
||||||
/** -------------------- Live Runs in Redis -------------------- */
|
/** -------------------- Live Runs in Redis -------------------- */
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user