Now all communication happens over sockets
This commit is contained in:
@@ -8,118 +8,28 @@ import { Cursor, cursorEffect } from "../../webapp/src/cursorEffect";
|
||||
import { Socket } from "socket.io";
|
||||
import { DiskStorage } from "./disk_storage";
|
||||
import { PageMeta } from "./server";
|
||||
import { Client, Page } from "./types";
|
||||
import { ClientPageState, Page } from "./types";
|
||||
import { safeRun } from "./util";
|
||||
|
||||
export class RealtimeStorage extends DiskStorage {
|
||||
export class SocketAPI {
|
||||
openPages = new Map<string, Page>();
|
||||
|
||||
private disconnectClient(client: Client, page: Page) {
|
||||
page.clients.delete(client);
|
||||
if (page.clients.size === 0) {
|
||||
console.log("No more clients for", page.name, "flushing");
|
||||
this.flushPageToDisk(page.name, page);
|
||||
this.openPages.delete(page.name);
|
||||
} else {
|
||||
page.cursors.delete(client.socket.id);
|
||||
this.broadcastCursors(page);
|
||||
}
|
||||
}
|
||||
|
||||
private broadcastCursors(page: Page) {
|
||||
page.clients.forEach((client) => {
|
||||
client.socket.emit("cursors", Object.fromEntries(page.cursors.entries()));
|
||||
});
|
||||
}
|
||||
|
||||
private flushPageToDisk(name: string, page: Page) {
|
||||
super
|
||||
.writePage(name, page.text.sliceString(0))
|
||||
.then((meta) => {
|
||||
console.log(`Wrote page ${name} to disk`);
|
||||
page.meta = meta;
|
||||
})
|
||||
.catch((e) => {
|
||||
console.log(`Could not write ${name} to disk:`, e);
|
||||
});
|
||||
}
|
||||
|
||||
// Override
|
||||
async readPage(pageName: string): Promise<{ text: string; meta: PageMeta }> {
|
||||
let page = this.openPages.get(pageName);
|
||||
if (page) {
|
||||
console.log("Serving page from memory", pageName);
|
||||
return {
|
||||
text: page.text.sliceString(0),
|
||||
meta: page.meta,
|
||||
};
|
||||
} else {
|
||||
return super.readPage(pageName);
|
||||
}
|
||||
}
|
||||
|
||||
async writePage(pageName: string, text: string): Promise<PageMeta> {
|
||||
let page = this.openPages.get(pageName);
|
||||
if (page) {
|
||||
for (let client of page.clients) {
|
||||
client.socket.emit("reload", pageName);
|
||||
}
|
||||
this.openPages.delete(pageName);
|
||||
}
|
||||
return super.writePage(pageName, text);
|
||||
}
|
||||
|
||||
disconnectPageSocket(socket: Socket, pageName: string) {
|
||||
let page = this.openPages.get(pageName);
|
||||
if (page) {
|
||||
for (let client of page.clients) {
|
||||
if (client.socket === socket) {
|
||||
this.disconnectClient(client, page);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
connectedSockets: Set<Socket> = new Set();
|
||||
pageStore: DiskStorage;
|
||||
|
||||
constructor(rootPath: string, io: Server) {
|
||||
super(rootPath);
|
||||
|
||||
// setInterval(() => {
|
||||
// console.log("Currently open pages:", this.openPages.keys());
|
||||
// }, 10000);
|
||||
|
||||
// Disk watcher
|
||||
fs.watch(
|
||||
rootPath,
|
||||
{
|
||||
recursive: true,
|
||||
persistent: false,
|
||||
},
|
||||
(eventType, filename) => {
|
||||
safeRun(async () => {
|
||||
if (path.extname(filename) !== ".md") {
|
||||
return;
|
||||
}
|
||||
let localPath = path.join(rootPath, filename);
|
||||
let pageName = filename.substring(0, filename.length - 3);
|
||||
let s = await stat(localPath);
|
||||
// console.log("Edit in", pageName);
|
||||
const openPage = this.openPages.get(pageName);
|
||||
if (openPage) {
|
||||
if (openPage.meta.lastModified < s.mtime.getTime()) {
|
||||
console.log("Page changed on disk outside of editor, reloading");
|
||||
this.openPages.delete(pageName);
|
||||
for (let client of openPage.clients) {
|
||||
client.socket.emit("reload", pageName);
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
);
|
||||
this.pageStore = new DiskStorage(rootPath);
|
||||
this.fileWatcher(rootPath);
|
||||
|
||||
io.on("connection", (socket) => {
|
||||
console.log("Connected", socket.id);
|
||||
let clientOpenPages = new Set<string>();
|
||||
this.connectedSockets.add(socket);
|
||||
const socketOpenPages = new Set<string>();
|
||||
|
||||
socket.on("disconnect", () => {
|
||||
console.log("Disconnected", socket.id);
|
||||
socketOpenPages.forEach(disconnectPageSocket);
|
||||
this.connectedSockets.delete(socket);
|
||||
});
|
||||
|
||||
function onCall(eventName: string, cb: (...args: any[]) => Promise<any>) {
|
||||
socket.on(eventName, (reqId: number, ...args) => {
|
||||
@@ -129,11 +39,23 @@ export class RealtimeStorage extends DiskStorage {
|
||||
});
|
||||
}
|
||||
|
||||
const _this = this;
|
||||
function disconnectPageSocket(pageName: string) {
|
||||
let page = _this.openPages.get(pageName);
|
||||
if (page) {
|
||||
for (let client of page.clientStates) {
|
||||
if (client.socket === socket) {
|
||||
_this.disconnectClient(client, page);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
onCall("openPage", async (pageName: string) => {
|
||||
let page = this.openPages.get(pageName);
|
||||
if (!page) {
|
||||
try {
|
||||
let { text, meta } = await super.readPage(pageName);
|
||||
let { text, meta } = await this.pageStore.readPage(pageName);
|
||||
page = new Page(pageName, text, meta);
|
||||
} catch (e) {
|
||||
console.log("Creating new page", pageName);
|
||||
@@ -141,8 +63,8 @@ export class RealtimeStorage extends DiskStorage {
|
||||
}
|
||||
this.openPages.set(pageName, page);
|
||||
}
|
||||
page.clients.add(new Client(socket, page.version));
|
||||
clientOpenPages.add(pageName);
|
||||
page.clientStates.add(new ClientPageState(socket, page.version));
|
||||
socketOpenPages.add(pageName);
|
||||
console.log("Opened page", pageName);
|
||||
this.broadcastCursors(page);
|
||||
return page.toJSON();
|
||||
@@ -150,8 +72,8 @@ export class RealtimeStorage extends DiskStorage {
|
||||
|
||||
socket.on("closePage", (pageName: string) => {
|
||||
console.log("Closing page", pageName);
|
||||
clientOpenPages.delete(pageName);
|
||||
this.disconnectPageSocket(socket, pageName);
|
||||
socketOpenPages.delete(pageName);
|
||||
disconnectPageSocket(pageName);
|
||||
});
|
||||
|
||||
onCall(
|
||||
@@ -169,13 +91,13 @@ export class RealtimeStorage extends DiskStorage {
|
||||
pageName,
|
||||
this.openPages.keys()
|
||||
);
|
||||
return;
|
||||
return false;
|
||||
}
|
||||
if (version !== page.version) {
|
||||
console.error("Invalid version", version, page.version);
|
||||
return false;
|
||||
} else {
|
||||
console.log("Applying", updates.length, "updates");
|
||||
console.log("Applying", updates.length, "updates to", pageName);
|
||||
let transformedUpdates = [];
|
||||
let textChanged = false;
|
||||
for (let update of updates) {
|
||||
@@ -225,7 +147,7 @@ export class RealtimeStorage extends DiskStorage {
|
||||
}
|
||||
// TODO: Optimize this
|
||||
let oldestVersion = Infinity;
|
||||
page.clients.forEach((client) => {
|
||||
page.clientStates.forEach((client) => {
|
||||
oldestVersion = Math.min(client.version, oldestVersion);
|
||||
if (client.socket === socket) {
|
||||
client.version = version;
|
||||
@@ -242,12 +164,139 @@ export class RealtimeStorage extends DiskStorage {
|
||||
}
|
||||
);
|
||||
|
||||
socket.on("disconnect", () => {
|
||||
console.log("Disconnected", socket.id);
|
||||
clientOpenPages.forEach((pageName) => {
|
||||
this.disconnectPageSocket(socket, pageName);
|
||||
});
|
||||
onCall(
|
||||
"readPage",
|
||||
async (pageName: string): Promise<{ text: string; meta: PageMeta }> => {
|
||||
let page = this.openPages.get(pageName);
|
||||
if (page) {
|
||||
console.log("Serving page from memory", pageName);
|
||||
return {
|
||||
text: page.text.sliceString(0),
|
||||
meta: page.meta,
|
||||
};
|
||||
} else {
|
||||
return this.pageStore.readPage(pageName);
|
||||
}
|
||||
}
|
||||
);
|
||||
|
||||
onCall("writePage", async (pageName: string, text: string) => {
|
||||
let page = this.openPages.get(pageName);
|
||||
if (page) {
|
||||
for (let client of page.clientStates) {
|
||||
client.socket.emit("reloadPage", pageName);
|
||||
}
|
||||
this.openPages.delete(pageName);
|
||||
}
|
||||
return this.pageStore.writePage(pageName, text);
|
||||
});
|
||||
|
||||
onCall("deletePage", async (pageName: string) => {
|
||||
this.openPages.delete(pageName);
|
||||
socketOpenPages.delete(pageName);
|
||||
// Cascading of this to all connected clients will be handled by file watcher
|
||||
return this.pageStore.deletePage(pageName);
|
||||
});
|
||||
|
||||
onCall("listPages", async (): Promise<PageMeta[]> => {
|
||||
return this.pageStore.listPages();
|
||||
});
|
||||
|
||||
onCall("getPageMeta", async (pageName: string): Promise<PageMeta> => {
|
||||
let page = this.openPages.get(pageName);
|
||||
if (page) {
|
||||
return page.meta;
|
||||
}
|
||||
return this.pageStore.getPageMeta(pageName);
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
private disconnectClient(client: ClientPageState, page: Page) {
|
||||
page.clientStates.delete(client);
|
||||
if (page.clientStates.size === 0) {
|
||||
console.log("No more clients for", page.name, "flushing");
|
||||
this.flushPageToDisk(page.name, page);
|
||||
this.openPages.delete(page.name);
|
||||
} else {
|
||||
page.cursors.delete(client.socket.id);
|
||||
this.broadcastCursors(page);
|
||||
}
|
||||
}
|
||||
|
||||
private broadcastCursors(page: Page) {
|
||||
page.clientStates.forEach((client) => {
|
||||
client.socket.emit(
|
||||
"cursorSnapshot",
|
||||
page.name,
|
||||
Object.fromEntries(page.cursors.entries())
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
private flushPageToDisk(name: string, page: Page) {
|
||||
safeRun(async () => {
|
||||
let meta = await this.pageStore.writePage(name, page.text.sliceString(0));
|
||||
console.log(`Wrote page ${name} to disk`);
|
||||
page.meta = meta;
|
||||
});
|
||||
}
|
||||
|
||||
private fileWatcher(rootPath: string) {
|
||||
fs.watch(
|
||||
rootPath,
|
||||
{
|
||||
recursive: true,
|
||||
persistent: false,
|
||||
},
|
||||
(eventType, filename) => {
|
||||
safeRun(async () => {
|
||||
if (path.extname(filename) !== ".md") {
|
||||
return;
|
||||
}
|
||||
let localPath = path.join(rootPath, filename);
|
||||
let pageName = filename.substring(0, filename.length - 3);
|
||||
// console.log("Edit in", pageName, eventType);
|
||||
let modifiedTime = 0;
|
||||
try {
|
||||
let s = await stat(localPath);
|
||||
modifiedTime = s.mtime.getTime();
|
||||
} catch (e) {
|
||||
// File was deleted
|
||||
console.log("Deleted", pageName);
|
||||
for (let socket of this.connectedSockets) {
|
||||
socket.emit("pageDeleted", pageName);
|
||||
}
|
||||
return;
|
||||
}
|
||||
const openPage = this.openPages.get(pageName);
|
||||
if (openPage) {
|
||||
if (openPage.meta.lastModified < modifiedTime) {
|
||||
console.log("Page changed on disk outside of editor, reloading");
|
||||
this.openPages.delete(pageName);
|
||||
const meta = {
|
||||
name: pageName,
|
||||
lastModified: modifiedTime,
|
||||
} as PageMeta;
|
||||
for (let client of openPage.clientStates) {
|
||||
client.socket.emit("pageChanged", meta);
|
||||
}
|
||||
}
|
||||
}
|
||||
if (eventType === "rename") {
|
||||
// This most likely means a new file was created, let's push new file listings to all connected sockets
|
||||
console.log(
|
||||
"New file created, broadcasting to all connected sockets"
|
||||
);
|
||||
for (let socket of this.connectedSockets) {
|
||||
socket.emit("pageCreated", {
|
||||
name: pageName,
|
||||
lastModified: modifiedTime,
|
||||
} as PageMeta);
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
);
|
||||
}
|
||||
}
|
||||
+2
-81
@@ -1,12 +1,8 @@
|
||||
import bodyParser from "body-parser";
|
||||
import cors from "cors";
|
||||
import express from "express";
|
||||
import { readFile } from "fs/promises";
|
||||
import http from "http";
|
||||
import { Server } from "socket.io";
|
||||
import stream from "stream";
|
||||
import { promisify } from "util";
|
||||
import { RealtimeStorage } from "./realtime_storage";
|
||||
import { SocketAPI } from "./api";
|
||||
|
||||
const app = express();
|
||||
const server = http.createServer(app);
|
||||
@@ -18,7 +14,6 @@ const io = new Server(server, {
|
||||
});
|
||||
|
||||
const port = 3000;
|
||||
const pipeline = promisify(stream.pipeline);
|
||||
export const pagesPath = "../pages";
|
||||
const distDir = `${__dirname}/../../webapp/dist`;
|
||||
|
||||
@@ -29,80 +24,7 @@ export type PageMeta = {
|
||||
};
|
||||
|
||||
app.use("/", express.static(distDir));
|
||||
|
||||
let fsRouter = express.Router();
|
||||
// let diskFS = new DiskFS(pagesPath);
|
||||
let filesystem = new RealtimeStorage(pagesPath, io);
|
||||
|
||||
// Page list
|
||||
fsRouter.route("/").get(async (req, res) => {
|
||||
res.json(await filesystem.listPages());
|
||||
});
|
||||
|
||||
fsRouter
|
||||
.route(/\/(.+)/)
|
||||
.get(async (req, res) => {
|
||||
let reqPath = req.params[0];
|
||||
console.log("Getting", reqPath);
|
||||
try {
|
||||
let { text, meta } = await filesystem.readPage(reqPath);
|
||||
res.status(200);
|
||||
res.header("Last-Modified", "" + meta.lastModified);
|
||||
res.header("Content-Type", "text/markdown");
|
||||
res.send(text);
|
||||
} catch (e) {
|
||||
res.status(200);
|
||||
res.send("");
|
||||
}
|
||||
})
|
||||
.put(bodyParser.text({ type: "*/*" }), async (req, res) => {
|
||||
let reqPath = req.params[0];
|
||||
|
||||
try {
|
||||
let meta = await filesystem.writePage(reqPath, req.body);
|
||||
res.status(200);
|
||||
res.header("Last-Modified", "" + meta.lastModified);
|
||||
res.send("OK");
|
||||
} catch (err) {
|
||||
res.status(500);
|
||||
res.send("Write failed");
|
||||
console.error("Pipeline failed", err);
|
||||
}
|
||||
})
|
||||
.options(async (req, res) => {
|
||||
let reqPath = req.params[0];
|
||||
try {
|
||||
const meta = await filesystem.getPageMeta(reqPath);
|
||||
res.status(200);
|
||||
res.header("Last-Modified", "" + meta.lastModified);
|
||||
res.header("Content-Type", "text/markdown");
|
||||
res.send("");
|
||||
} catch (e) {
|
||||
res.status(200);
|
||||
res.send("");
|
||||
}
|
||||
})
|
||||
.delete(async (req, res) => {
|
||||
let reqPath = req.params[0];
|
||||
try {
|
||||
await filesystem.deletePage(reqPath);
|
||||
res.status(200);
|
||||
res.send("OK");
|
||||
} catch (e) {
|
||||
console.error("Error deleting file", reqPath, e);
|
||||
res.status(500);
|
||||
res.send("OK");
|
||||
}
|
||||
});
|
||||
|
||||
app.use(
|
||||
"/fs",
|
||||
cors({
|
||||
methods: "GET,HEAD,PUT,OPTIONS,POST,DELETE",
|
||||
preflightContinue: true,
|
||||
}),
|
||||
fsRouter
|
||||
);
|
||||
let filesystem = new SocketAPI(pagesPath, io);
|
||||
|
||||
// Fallback, serve index.html
|
||||
let cachedIndex: string | undefined = undefined;
|
||||
@@ -113,7 +35,6 @@ app.get("/*", async (req, res) => {
|
||||
res.status(200).header("Content-Type", "text/html").send(cachedIndex);
|
||||
});
|
||||
|
||||
//sup
|
||||
server.listen(port, () => {
|
||||
console.log(`Server istening on port ${port}`);
|
||||
});
|
||||
|
||||
+2
-2
@@ -4,7 +4,7 @@ import { Socket } from "socket.io";
|
||||
import { Cursor } from "../../webapp/src/cursorEffect";
|
||||
import { PageMeta } from "./server";
|
||||
|
||||
export class Client {
|
||||
export class ClientPageState {
|
||||
constructor(public socket: Socket, public version: number) {}
|
||||
}
|
||||
|
||||
@@ -12,7 +12,7 @@ export class Page {
|
||||
versionOffset = 0;
|
||||
updates: Update[] = [];
|
||||
cursors = new Map<string, Cursor>();
|
||||
clients = new Set<Client>();
|
||||
clientStates = new Set<ClientPageState>();
|
||||
|
||||
pending: ((value: any) => void)[] = [];
|
||||
|
||||
|
||||
Reference in New Issue
Block a user