Merge pull request #3 from orus-dev/full-websocket-system
Full websocket system
This commit is contained in:
Generated
+33
-1
@@ -11,10 +11,12 @@
|
|||||||
"@types/cors": "^2.8.19",
|
"@types/cors": "^2.8.19",
|
||||||
"@types/diff-match-patch": "^1.0.36",
|
"@types/diff-match-patch": "^1.0.36",
|
||||||
"@types/express": "^5.0.6",
|
"@types/express": "^5.0.6",
|
||||||
|
"@types/ws": "^8.18.1",
|
||||||
"axios": "^1.13.5",
|
"axios": "^1.13.5",
|
||||||
"cors": "^2.8.6",
|
"cors": "^2.8.6",
|
||||||
"diff-match-patch": "^1.0.5",
|
"diff-match-patch": "^1.0.5",
|
||||||
"express": "^5.2.1"
|
"express": "^5.2.1",
|
||||||
|
"ws": "^8.19.0"
|
||||||
},
|
},
|
||||||
"devDependencies": {
|
"devDependencies": {
|
||||||
"@types/mocha": "^10.0.6",
|
"@types/mocha": "^10.0.6",
|
||||||
@@ -407,6 +409,15 @@
|
|||||||
"dev": true,
|
"dev": true,
|
||||||
"license": "MIT"
|
"license": "MIT"
|
||||||
},
|
},
|
||||||
|
"node_modules/@types/ws": {
|
||||||
|
"version": "8.18.1",
|
||||||
|
"resolved": "https://registry.npmjs.org/@types/ws/-/ws-8.18.1.tgz",
|
||||||
|
"integrity": "sha512-ThVF6DCVhA8kUGy+aazFQ4kXQ7E1Ty7A3ypFOe0IcJV8O/M511G99AW24irKrW56Wt44yG9+ij8FaqoBGkuBXg==",
|
||||||
|
"license": "MIT",
|
||||||
|
"dependencies": {
|
||||||
|
"@types/node": "*"
|
||||||
|
}
|
||||||
|
},
|
||||||
"node_modules/@typescript-eslint/eslint-plugin": {
|
"node_modules/@typescript-eslint/eslint-plugin": {
|
||||||
"version": "6.21.0",
|
"version": "6.21.0",
|
||||||
"resolved": "https://registry.npmjs.org/@typescript-eslint/eslint-plugin/-/eslint-plugin-6.21.0.tgz",
|
"resolved": "https://registry.npmjs.org/@typescript-eslint/eslint-plugin/-/eslint-plugin-6.21.0.tgz",
|
||||||
@@ -4140,6 +4151,27 @@
|
|||||||
"integrity": "sha512-l4Sp/DRseor9wL6EvV2+TuQn63dMkPjZ/sp9XkghTEbV9KlPS1xUsZ3u7/IQO4wxtcFB4bgpQPRcR3QCvezPcQ==",
|
"integrity": "sha512-l4Sp/DRseor9wL6EvV2+TuQn63dMkPjZ/sp9XkghTEbV9KlPS1xUsZ3u7/IQO4wxtcFB4bgpQPRcR3QCvezPcQ==",
|
||||||
"license": "ISC"
|
"license": "ISC"
|
||||||
},
|
},
|
||||||
|
"node_modules/ws": {
|
||||||
|
"version": "8.19.0",
|
||||||
|
"resolved": "https://registry.npmjs.org/ws/-/ws-8.19.0.tgz",
|
||||||
|
"integrity": "sha512-blAT2mjOEIi0ZzruJfIhb3nps74PRWTCz1IjglWEEpQl5XS/UNama6u2/rjFkDDouqr4L67ry+1aGIALViWjDg==",
|
||||||
|
"license": "MIT",
|
||||||
|
"engines": {
|
||||||
|
"node": ">=10.0.0"
|
||||||
|
},
|
||||||
|
"peerDependencies": {
|
||||||
|
"bufferutil": "^4.0.1",
|
||||||
|
"utf-8-validate": ">=5.0.2"
|
||||||
|
},
|
||||||
|
"peerDependenciesMeta": {
|
||||||
|
"bufferutil": {
|
||||||
|
"optional": true
|
||||||
|
},
|
||||||
|
"utf-8-validate": {
|
||||||
|
"optional": true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
},
|
||||||
"node_modules/y18n": {
|
"node_modules/y18n": {
|
||||||
"version": "5.0.8",
|
"version": "5.0.8",
|
||||||
"resolved": "https://registry.npmjs.org/y18n/-/y18n-5.0.8.tgz",
|
"resolved": "https://registry.npmjs.org/y18n/-/y18n-5.0.8.tgz",
|
||||||
|
|||||||
+3
-1
@@ -91,9 +91,11 @@
|
|||||||
"@types/cors": "^2.8.19",
|
"@types/cors": "^2.8.19",
|
||||||
"@types/diff-match-patch": "^1.0.36",
|
"@types/diff-match-patch": "^1.0.36",
|
||||||
"@types/express": "^5.0.6",
|
"@types/express": "^5.0.6",
|
||||||
|
"@types/ws": "^8.18.1",
|
||||||
"axios": "^1.13.5",
|
"axios": "^1.13.5",
|
||||||
"cors": "^2.8.6",
|
"cors": "^2.8.6",
|
||||||
"diff-match-patch": "^1.0.5",
|
"diff-match-patch": "^1.0.5",
|
||||||
"express": "^5.2.1"
|
"express": "^5.2.1",
|
||||||
|
"ws": "^8.19.0"
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+187
-45
@@ -1,51 +1,200 @@
|
|||||||
import axios from "axios";
|
import { WebSocket } from "ws";
|
||||||
import { LiveRunMove } from "./types";
|
import { LiveRunMove } from "./types";
|
||||||
|
import { randomUUID } from "crypto";
|
||||||
|
|
||||||
var cookies: string | undefined;
|
/* ──────────────────────────
|
||||||
|
State
|
||||||
|
────────────────────────── */
|
||||||
|
|
||||||
export function getCookies(): string | undefined {
|
let socket: WebSocket | null = null;
|
||||||
|
let connecting: Promise<WebSocket> | null = null;
|
||||||
|
let currentOrigin: string | null = null;
|
||||||
|
|
||||||
|
let cookies: Record<string, string> | undefined;
|
||||||
|
|
||||||
|
type PendingResolver = {
|
||||||
|
resolve: (value: any) => void;
|
||||||
|
reject: (err: any) => void;
|
||||||
|
timeout: NodeJS.Timeout;
|
||||||
|
};
|
||||||
|
|
||||||
|
const pending = new Map<string, PendingResolver>();
|
||||||
|
|
||||||
|
/* ──────────────────────────
|
||||||
|
Cookies
|
||||||
|
────────────────────────── */
|
||||||
|
|
||||||
|
export function getCookies(): Record<string, string> | undefined {
|
||||||
return cookies;
|
return cookies;
|
||||||
}
|
}
|
||||||
export function setCookies(c: string) {
|
|
||||||
|
export function setCookies(c: Record<string, string>) {
|
||||||
cookies = c;
|
cookies = c;
|
||||||
}
|
}
|
||||||
function getOrigin(useLocalhost: boolean) {
|
|
||||||
return useLocalhost ? "http://localhost:3000" : "https://dev-run.selimaj.dev";
|
/* ──────────────────────────
|
||||||
|
Utils
|
||||||
|
────────────────────────── */
|
||||||
|
|
||||||
|
function getWsOrigin(useLocalhost: boolean) {
|
||||||
|
return useLocalhost
|
||||||
|
? "ws://localhost:3000/api/run"
|
||||||
|
: "wss://dev-run.selimaj.dev/api/run";
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function rejectAllPending(err: Error) {
|
||||||
|
for (const { reject, timeout } of pending.values()) {
|
||||||
|
clearTimeout(timeout);
|
||||||
|
reject(err);
|
||||||
|
}
|
||||||
|
pending.clear();
|
||||||
|
}
|
||||||
|
|
||||||
|
/* ──────────────────────────
|
||||||
|
Socket management
|
||||||
|
────────────────────────── */
|
||||||
|
|
||||||
|
async function ensureSocket(useLocalhost: boolean): Promise<WebSocket> {
|
||||||
|
const origin = getWsOrigin(useLocalhost);
|
||||||
|
|
||||||
|
if (
|
||||||
|
socket &&
|
||||||
|
socket.readyState === WebSocket.OPEN &&
|
||||||
|
currentOrigin === origin
|
||||||
|
) {
|
||||||
|
return socket;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (connecting) {
|
||||||
|
return connecting;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (socket && currentOrigin !== origin) {
|
||||||
|
socket.close();
|
||||||
|
socket = null;
|
||||||
|
}
|
||||||
|
|
||||||
|
connecting = new Promise<WebSocket>((resolve, reject) => {
|
||||||
|
const ws = new WebSocket(origin);
|
||||||
|
|
||||||
|
ws.onopen = async () => {
|
||||||
|
try {
|
||||||
|
const authCookies = getCookies();
|
||||||
|
if (!authCookies) {
|
||||||
|
throw new Error("Cookies not set");
|
||||||
|
}
|
||||||
|
|
||||||
|
socket = ws;
|
||||||
|
currentOrigin = origin;
|
||||||
|
connecting = null;
|
||||||
|
|
||||||
|
await sendUnsafe(useLocalhost, {
|
||||||
|
type: "auth",
|
||||||
|
cookies: authCookies,
|
||||||
|
});
|
||||||
|
|
||||||
|
resolve(ws);
|
||||||
|
} catch (err) {
|
||||||
|
ws.close();
|
||||||
|
reject(err);
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
ws.onmessage = (e) => {
|
||||||
|
let msg: any;
|
||||||
|
try {
|
||||||
|
msg = JSON.parse(e.data.toString());
|
||||||
|
} catch {
|
||||||
|
console.warn("Invalid WS message:", e.data);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
const { requestId, ok, data, error } = msg;
|
||||||
|
if (!requestId) return;
|
||||||
|
|
||||||
|
const pendingReq = pending.get(requestId);
|
||||||
|
if (!pendingReq) return;
|
||||||
|
|
||||||
|
pending.delete(requestId);
|
||||||
|
clearTimeout(pendingReq.timeout);
|
||||||
|
|
||||||
|
ok ? pendingReq.resolve(data) : pendingReq.reject(error);
|
||||||
|
};
|
||||||
|
|
||||||
|
ws.onerror = (err) => {
|
||||||
|
rejectAllPending(new Error("WebSocket error"));
|
||||||
|
connecting = null;
|
||||||
|
reject(err);
|
||||||
|
};
|
||||||
|
|
||||||
|
ws.onclose = () => {
|
||||||
|
rejectAllPending(new Error("WebSocket closed"));
|
||||||
|
socket = null;
|
||||||
|
currentOrigin = null;
|
||||||
|
connecting = null;
|
||||||
|
};
|
||||||
|
});
|
||||||
|
|
||||||
|
return connecting;
|
||||||
|
}
|
||||||
|
|
||||||
|
/* ──────────────────────────
|
||||||
|
RPC helpers
|
||||||
|
────────────────────────── */
|
||||||
|
|
||||||
|
async function sendUnsafe<T>(
|
||||||
|
useLocalhost: boolean,
|
||||||
|
payload: Record<string, any>,
|
||||||
|
): Promise<T> {
|
||||||
|
if (!socket || socket.readyState !== WebSocket.OPEN) {
|
||||||
|
throw new Error("WebSocket is not open");
|
||||||
|
}
|
||||||
|
|
||||||
|
const requestId = randomUUID();
|
||||||
|
|
||||||
|
return new Promise<T>((resolve, reject) => {
|
||||||
|
const timeout = setTimeout(() => {
|
||||||
|
pending.delete(requestId);
|
||||||
|
reject(new Error("WebSocket request timed out"));
|
||||||
|
}, 10_000);
|
||||||
|
|
||||||
|
pending.set(requestId, { resolve, reject, timeout });
|
||||||
|
|
||||||
|
socket?.send(JSON.stringify({ ...payload, requestId }));
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
async function send<T>(
|
||||||
|
useLocalhost: boolean,
|
||||||
|
payload: Record<string, any>,
|
||||||
|
): Promise<T> {
|
||||||
|
await ensureSocket(useLocalhost);
|
||||||
|
return sendUnsafe(useLocalhost, payload);
|
||||||
|
}
|
||||||
|
|
||||||
|
/* ──────────────────────────
|
||||||
|
Public API
|
||||||
|
────────────────────────── */
|
||||||
|
|
||||||
export async function addRun(
|
export async function addRun(
|
||||||
useLocalhost: boolean,
|
useLocalhost: boolean,
|
||||||
problem: string,
|
problem: string,
|
||||||
mode: string,
|
mode: string,
|
||||||
): Promise<string> {
|
): Promise<string> {
|
||||||
return (
|
const data = await send<{ runId: string }>(useLocalhost, {
|
||||||
await axios.put(
|
type: "create",
|
||||||
getOrigin(useLocalhost) + "/api/run",
|
problem,
|
||||||
{
|
category: mode,
|
||||||
problem,
|
});
|
||||||
category: mode,
|
|
||||||
},
|
return data.runId;
|
||||||
{
|
|
||||||
headers: {
|
|
||||||
Cookie: cookies,
|
|
||||||
},
|
|
||||||
},
|
|
||||||
)
|
|
||||||
).data.runId;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
export async function submitRun(useLocalhost: boolean, runId: string) {
|
export async function submitRun(useLocalhost: boolean, runId: string) {
|
||||||
await axios.post(
|
await send(useLocalhost, {
|
||||||
getOrigin(useLocalhost) + "/api/run",
|
type: "submit",
|
||||||
{
|
runId,
|
||||||
runId,
|
});
|
||||||
},
|
|
||||||
{
|
|
||||||
headers: {
|
|
||||||
Cookie: cookies,
|
|
||||||
},
|
|
||||||
},
|
|
||||||
);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
export async function addRunMoves(
|
export async function addRunMoves(
|
||||||
@@ -55,18 +204,11 @@ export async function addRunMoves(
|
|||||||
language: string | null,
|
language: string | null,
|
||||||
moves: LiveRunMove[],
|
moves: LiveRunMove[],
|
||||||
) {
|
) {
|
||||||
await axios.post(
|
await send(useLocalhost, {
|
||||||
getOrigin(useLocalhost) + "/api/run/move",
|
type: "move",
|
||||||
{
|
runId,
|
||||||
runId,
|
file,
|
||||||
file,
|
language,
|
||||||
moves,
|
moves,
|
||||||
language,
|
});
|
||||||
},
|
|
||||||
{
|
|
||||||
headers: {
|
|
||||||
Cookie: cookies,
|
|
||||||
},
|
|
||||||
},
|
|
||||||
);
|
|
||||||
}
|
}
|
||||||
|
|||||||
+56
-21
@@ -11,6 +11,7 @@ export function startMonitoring(useLocalhost: boolean, runId: string) {
|
|||||||
let moves: LiveRunMove[] = [];
|
let moves: LiveRunMove[] = [];
|
||||||
|
|
||||||
let lastText = "";
|
let lastText = "";
|
||||||
|
let lastCursorOffset: number | undefined;
|
||||||
let lastEventTime = Date.now();
|
let lastEventTime = Date.now();
|
||||||
let idleTimeout: NodeJS.Timeout | undefined;
|
let idleTimeout: NodeJS.Timeout | undefined;
|
||||||
let isSending = false;
|
let isSending = false;
|
||||||
@@ -18,7 +19,9 @@ export function startMonitoring(useLocalhost: boolean, runId: string) {
|
|||||||
const IDLE_MS = 350;
|
const IDLE_MS = 350;
|
||||||
|
|
||||||
if (vscode.window.activeTextEditor) {
|
if (vscode.window.activeTextEditor) {
|
||||||
lastText = vscode.window.activeTextEditor.document.getText();
|
const editor = vscode.window.activeTextEditor;
|
||||||
|
lastText = editor.document.getText();
|
||||||
|
lastCursorOffset = editor.document.offsetAt(editor.selection.active);
|
||||||
}
|
}
|
||||||
|
|
||||||
async function flush(file: string | null, language: string | null) {
|
async function flush(file: string | null, language: string | null) {
|
||||||
@@ -42,41 +45,72 @@ export function startMonitoring(useLocalhost: boolean, runId: string) {
|
|||||||
return latency;
|
return latency;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// --- Existing text change interval ---
|
||||||
const readInterval = setInterval(() => {
|
const readInterval = setInterval(() => {
|
||||||
const editor = vscode.window.activeTextEditor;
|
const editor = vscode.window.activeTextEditor;
|
||||||
if (!editor) return;
|
if (!editor) return;
|
||||||
|
|
||||||
const newText = editor.document.getText();
|
const newText = editor.document.getText();
|
||||||
if (newText === lastText) return;
|
if (newText !== lastText) {
|
||||||
|
let latency = getLatency();
|
||||||
|
|
||||||
let latency = getLatency();
|
const diffs = dmp.diff_main(lastText, newText);
|
||||||
|
dmp.diff_cleanupEfficiency(diffs);
|
||||||
|
|
||||||
const diffs = dmp.diff_main(lastText, newText);
|
const changes = diffsToChanges(diffs);
|
||||||
dmp.diff_cleanupEfficiency(diffs);
|
|
||||||
|
|
||||||
const changes = diffsToChanges(diffs);
|
lastText = newText;
|
||||||
|
|
||||||
lastText = newText;
|
for (const change of changes) {
|
||||||
|
moves.push({
|
||||||
|
moveId: moveId++,
|
||||||
|
latency,
|
||||||
|
cursor: editor.document.offsetAt(editor.selection.active),
|
||||||
|
changes: change,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
if (idleTimeout) clearTimeout(idleTimeout);
|
||||||
|
idleTimeout = setTimeout(
|
||||||
|
() =>
|
||||||
|
flush(
|
||||||
|
path.basename(editor?.document.fileName || "") || null,
|
||||||
|
editor?.document.languageId || null,
|
||||||
|
),
|
||||||
|
IDLE_MS,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}, 500);
|
||||||
|
|
||||||
|
// --- New cursor tracking interval ---
|
||||||
|
const cursorInterval = setInterval(() => {
|
||||||
|
const editor = vscode.window.activeTextEditor;
|
||||||
|
if (!editor) return;
|
||||||
|
|
||||||
|
const cursorOffset = editor.document.offsetAt(editor.selection.active);
|
||||||
|
|
||||||
|
if (cursorOffset !== lastCursorOffset) {
|
||||||
|
const latency = getLatency();
|
||||||
|
|
||||||
for (const change of changes) {
|
|
||||||
moves.push({
|
moves.push({
|
||||||
moveId: moveId++,
|
moveId: moveId++,
|
||||||
latency,
|
latency,
|
||||||
cursor: editor.document.offsetAt(editor.selection.active),
|
cursor: cursorOffset,
|
||||||
changes: change,
|
|
||||||
});
|
});
|
||||||
}
|
|
||||||
|
|
||||||
if (idleTimeout) clearTimeout(idleTimeout);
|
lastCursorOffset = cursorOffset;
|
||||||
idleTimeout = setTimeout(
|
|
||||||
() =>
|
if (idleTimeout) clearTimeout(idleTimeout);
|
||||||
flush(
|
idleTimeout = setTimeout(
|
||||||
path.basename(editor?.document.fileName || "") || null,
|
() =>
|
||||||
editor?.document.languageId || null,
|
flush(
|
||||||
),
|
path.basename(editor.document.fileName || "") || null,
|
||||||
IDLE_MS,
|
editor.document.languageId || null,
|
||||||
);
|
),
|
||||||
}, 500);
|
IDLE_MS,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}, 1000); // every second
|
||||||
|
|
||||||
vscode.window.onDidChangeActiveTextEditor((editor) => {
|
vscode.window.onDidChangeActiveTextEditor((editor) => {
|
||||||
let latency = getLatency();
|
let latency = getLatency();
|
||||||
@@ -99,6 +133,7 @@ export function startMonitoring(useLocalhost: boolean, runId: string) {
|
|||||||
|
|
||||||
return () => {
|
return () => {
|
||||||
clearInterval(readInterval);
|
clearInterval(readInterval);
|
||||||
|
clearInterval(cursorInterval);
|
||||||
if (idleTimeout) clearTimeout(idleTimeout);
|
if (idleTimeout) clearTimeout(idleTimeout);
|
||||||
flush(
|
flush(
|
||||||
path.basename(vscode.window.activeTextEditor?.document.fileName || "") ||
|
path.basename(vscode.window.activeTextEditor?.document.fileName || "") ||
|
||||||
|
|||||||
Reference in New Issue
Block a user