Working websocket structure

This commit is contained in:
2026-02-09 14:00:00 +01:00
parent 6c9b703ed1
commit 33f467cf02
4 changed files with 66 additions and 14 deletions
+11 -1
View File
@@ -19,6 +19,15 @@ export function UPGRADE(client: WebSocket, server: WebSocketServer) {
client.once("message", (message) => {
runId = message.toString();
const sendText = async () => {
const text = await subscriber.get(`liveRunText:${runId}`);
client.send(JSON.stringify({ moves: [], text }));
};
sendText();
client.on("message", sendText);
subscriber.subscribe(`liveRunMoves:${runId}`, (err, count) => {
if (err) {
console.error("Failed to subscribe", err);
@@ -29,8 +38,9 @@ export function UPGRADE(client: WebSocket, server: WebSocketServer) {
});
subscriber.on("message", (channel, message) => {
const moves = JSON.parse(message);
if (client.readyState === WebSocket.OPEN) {
client.send(message);
client.send(JSON.stringify({ moves, text: null }));
}
});
});
+23 -9
View File
@@ -42,6 +42,7 @@ function getLanguageExtension(lang: Language) {
export default function Editor({ run }: { run: LiveRun | null | undefined }) {
const [editorView, setEditorView] = useState<EditorView>();
const scheduledTimeouts = useRef<NodeJS.Timeout[]>([]);
const [text, setText] = useState("");
// WebSocket connection
useEffect(() => {
@@ -54,17 +55,30 @@ export default function Editor({ run }: { run: LiveRun | null | undefined }) {
};
ws.onmessage = (m) => {
const data: LiveRunMove[] = JSON.parse(m.data);
const data: { moves: LiveRunMove[]; text: string | null } = JSON.parse(
m.data,
);
data.forEach((move) => {
if (data.text) {
setText(data.text);
}
let cancelled = false;
data.moves.forEach((move) => {
const timeout = setTimeout(() => {
if (!editorView) return;
if (!editorView || cancelled) return;
editorView.dispatch({
selection: { anchor: move.cursor },
scrollIntoView: true,
changes: move.changes,
});
try {
editorView.dispatch({
selection: { anchor: move.cursor },
scrollIntoView: true,
changes: move.changes,
});
} catch {
cancelled = true;
ws.send("getText");
}
}, move.latency);
scheduledTimeouts.current.push(timeout);
@@ -78,7 +92,7 @@ export default function Editor({ run }: { run: LiveRun | null | undefined }) {
<CodeMirror
readOnly
className="h-full w-full"
value={``}
value={text}
extensions={[
getLanguageExtension("javascript"),
EditorState.transactionFilter.of((tr) => {
+16
View File
@@ -0,0 +1,16 @@
import { LiveRunMove } from "@/modules/live-run/types";
export function applyMoves(text: string, moves: LiveRunMove[]) {
let result = text;
for (const move of moves) {
if (!move.changes) continue;
result =
result.slice(0, move.changes.from) +
move.changes.insert +
result.slice(move.changes.to);
}
return result;
}
+16 -4
View File
@@ -5,6 +5,7 @@ import { LiveRun, LiveRunMove } from "./types";
import { SupabaseClient } from "@supabase/supabase-js";
import { redis } from "@/lib/redis";
import Redis from "ioredis";
import { applyMoves } from "@/lib/move";
/** -------------------- Live Runs in Redis -------------------- */
@@ -50,14 +51,25 @@ export async function removeLiveRun(id: string) {
* Add moves to a live run (ephemeral)
*/
export async function addLiveRunMoves(runId: string, moves: LiveRunMove[]) {
const runExists = await redis.exists(`liveRun:${runId}`);
const redis = new Redis();
// Make sure the live run exists
const runExists = await redis.exists(`liveRun:${runId}`);
if (!runExists) throw "Invalid live run";
const publisher = new Redis();
// Publish the new moves to subscribers
await redis.publish(`liveRunMoves:${runId}`, JSON.stringify(moves));
// Send a message
await publisher.publish(`liveRunMoves:${runId}`, JSON.stringify(moves));
// Get the last text (from Redis)
const lastText = (await redis.get(`liveRunText:${runId}`)) || "";
// Apply the new moves to get the updated text
const updatedText = applyMoves(lastText, moves);
// Store the updated text back in Redis
await redis.set(`liveRunText:${runId}`, updatedText);
return updatedText;
}
/** -------------------- Submit Run -------------------- */