Removed all traces of sockets, real-time collab and other stuff.
This commit is contained in:
@@ -1,102 +0,0 @@
|
||||
import { afterAll, beforeAll, describe, expect, test } from "@jest/globals";
|
||||
|
||||
import { createServer } from "http";
|
||||
import { io as Client } from "socket.io-client";
|
||||
import { Server } from "socket.io";
|
||||
import { SocketServer } from "./api_server";
|
||||
import * as path from "path";
|
||||
import * as fs from "fs";
|
||||
import { SilverBulletHooks } from "../common/manifest";
|
||||
import { System } from "../plugos/system";
|
||||
|
||||
describe("Server test", () => {
|
||||
let io: Server,
|
||||
socketServer: SocketServer,
|
||||
clientSocket: any,
|
||||
reqId = 0;
|
||||
const tmpDir = path.join(__dirname, "test");
|
||||
|
||||
function wsCall(eventName: string, ...args: any[]): Promise<any> {
|
||||
return new Promise((resolve, reject) => {
|
||||
reqId++;
|
||||
clientSocket.once(`${eventName}Resp${reqId}`, (err: any, result: any) => {
|
||||
if (err) {
|
||||
reject(err);
|
||||
} else {
|
||||
resolve(result);
|
||||
}
|
||||
});
|
||||
clientSocket.emit(eventName, reqId, ...args);
|
||||
});
|
||||
}
|
||||
|
||||
beforeAll((done) => {
|
||||
const httpServer = createServer();
|
||||
io = new Server(httpServer);
|
||||
fs.mkdirSync(tmpDir, { recursive: true });
|
||||
fs.writeFileSync(`${tmpDir}/test.md`, "This is a simple test");
|
||||
httpServer.listen(async () => {
|
||||
// @ts-ignore
|
||||
const port = httpServer.address().port;
|
||||
// @ts-ignore
|
||||
clientSocket = new Client(`http://localhost:${port}`);
|
||||
socketServer = new SocketServer(
|
||||
tmpDir,
|
||||
io,
|
||||
new System<SilverBulletHooks>("server")
|
||||
);
|
||||
clientSocket.on("connect", done);
|
||||
await socketServer.init();
|
||||
});
|
||||
});
|
||||
|
||||
afterAll(() => {
|
||||
io.close();
|
||||
clientSocket.close();
|
||||
socketServer.close();
|
||||
fs.rmSync(tmpDir, { recursive: true, force: true });
|
||||
});
|
||||
|
||||
test("List pages", async () => {
|
||||
let pages = await wsCall("page.listPages");
|
||||
expect(pages.length).toBe(1);
|
||||
await wsCall("page.writePage", "test2.md", "This is another test");
|
||||
let pages2 = await wsCall("page.listPages");
|
||||
expect(pages2.length).toBe(2);
|
||||
await wsCall("page.deletePage", "test2.md");
|
||||
let pages3 = await wsCall("page.listPages");
|
||||
expect(pages3.length).toBe(1);
|
||||
});
|
||||
|
||||
test("Index operations", async () => {
|
||||
await wsCall("index.clearPageIndexForPage", "test");
|
||||
await wsCall("index.set", "test", "testkey", "value");
|
||||
expect(await wsCall("index.get", "test", "testkey")).toBe("value");
|
||||
await wsCall("index.delete", "test", "testkey");
|
||||
expect(await wsCall("index.get", "test", "testkey")).toBe(null);
|
||||
await wsCall("index.set", "test", "unrelated", 10);
|
||||
await wsCall("index.set", "test", "unrelated", 12);
|
||||
await wsCall("index.set", "test2", "complicated", {
|
||||
name: "Bla",
|
||||
age: 123123,
|
||||
});
|
||||
await wsCall("index.set", "test", "complicated", { name: "Bla", age: 100 });
|
||||
await wsCall("index.set", "test", "complicated2", {
|
||||
name: "Bla",
|
||||
age: 101,
|
||||
});
|
||||
expect(await wsCall("index.get", "test", "complicated")).toStrictEqual({
|
||||
name: "Bla",
|
||||
age: 100,
|
||||
});
|
||||
let result = await wsCall("index.scanPrefixForPage", "test", "compli");
|
||||
expect(result.length).toBe(2);
|
||||
let result2 = await wsCall("index.scanPrefixGlobal", "compli");
|
||||
expect(result2.length).toBe(3);
|
||||
await wsCall("index.deletePrefixForPage", "test", "compli");
|
||||
let result3 = await wsCall("index.scanPrefixForPage", "test", "compli");
|
||||
expect(result3.length).toBe(0);
|
||||
let result4 = await wsCall("index.scanPrefixGlobal", "compli");
|
||||
expect(result4.length).toBe(1);
|
||||
});
|
||||
});
|
||||
@@ -1,150 +0,0 @@
|
||||
import { Server, Socket } from "socket.io";
|
||||
import { Page } from "./types";
|
||||
import * as path from "path";
|
||||
import { IndexApi } from "./index_api";
|
||||
import { PageApi } from "./page_api";
|
||||
import { SilverBulletHooks } from "../common/manifest";
|
||||
import { pageIndexSyscalls } from "./syscalls/page_index";
|
||||
import { safeRun } from "./util";
|
||||
import { System } from "../plugos/system";
|
||||
|
||||
export class ClientConnection {
|
||||
openPages = new Set<string>();
|
||||
|
||||
constructor(readonly sock: Socket) {}
|
||||
}
|
||||
|
||||
export interface ApiProvider {
|
||||
init(): Promise<void>;
|
||||
|
||||
api(): Object;
|
||||
}
|
||||
|
||||
export class SocketServer {
|
||||
private openPages = new Map<string, Page>();
|
||||
private connectedSockets = new Set<Socket>();
|
||||
private apis = new Map<string, ApiProvider>();
|
||||
readonly rootPath: string;
|
||||
private serverSocket: Server;
|
||||
system: System<SilverBulletHooks>;
|
||||
|
||||
constructor(
|
||||
rootPath: string,
|
||||
serverSocket: Server,
|
||||
system: System<SilverBulletHooks>
|
||||
) {
|
||||
this.rootPath = path.resolve(rootPath);
|
||||
this.serverSocket = serverSocket;
|
||||
this.system = system;
|
||||
}
|
||||
|
||||
async registerApi(name: string, apiProvider: ApiProvider) {
|
||||
await apiProvider.init();
|
||||
this.apis.set(name, apiProvider);
|
||||
}
|
||||
|
||||
public async init() {
|
||||
const indexApi = new IndexApi(this.rootPath);
|
||||
await this.registerApi("index", indexApi);
|
||||
this.system.registerSyscalls("indexer", [], pageIndexSyscalls(indexApi.db));
|
||||
await this.registerApi(
|
||||
"page",
|
||||
new PageApi(
|
||||
this.rootPath,
|
||||
this.connectedSockets,
|
||||
this.openPages,
|
||||
this.system
|
||||
)
|
||||
);
|
||||
|
||||
this.serverSocket.on("connection", (socket) => {
|
||||
const clientConn = new ClientConnection(socket);
|
||||
|
||||
console.log("Connected", socket.id);
|
||||
this.connectedSockets.add(socket);
|
||||
|
||||
socket.on("disconnect", () => {
|
||||
console.log("Disconnected", socket.id);
|
||||
clientConn.openPages.forEach((pageName) => {
|
||||
safeRun(async () => {
|
||||
await disconnectPageSocket(pageName);
|
||||
});
|
||||
});
|
||||
this.connectedSockets.delete(socket);
|
||||
});
|
||||
|
||||
socket.on("page.closePage", (pageName: string) => {
|
||||
console.log("Client closed page", pageName);
|
||||
safeRun(async () => {
|
||||
await disconnectPageSocket(pageName);
|
||||
});
|
||||
clientConn.openPages.delete(pageName);
|
||||
});
|
||||
|
||||
const onCall = (
|
||||
eventName: string,
|
||||
cb: (...args: any[]) => Promise<any>
|
||||
) => {
|
||||
socket.on(eventName, (reqId: number, ...args) => {
|
||||
cb(...args)
|
||||
.then((result) => {
|
||||
socket.emit(`${eventName}Resp${reqId}`, null, result);
|
||||
})
|
||||
.catch((err) => {
|
||||
socket.emit(`${eventName}Resp${reqId}`, err.message);
|
||||
});
|
||||
});
|
||||
};
|
||||
|
||||
const disconnectPageSocket = async (pageName: string) => {
|
||||
let page = this.openPages.get(pageName);
|
||||
if (page) {
|
||||
for (let client of page.clientStates) {
|
||||
if (client.socket === socket) {
|
||||
await (this.apis.get("page")! as PageApi).disconnectClient(
|
||||
client,
|
||||
page
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
};
|
||||
for (let [apiName, apiProvider] of this.apis) {
|
||||
Object.entries(apiProvider.api()).forEach(([eventName, cb]) => {
|
||||
onCall(`${apiName}.${eventName}`, (...args: any[]): any => {
|
||||
// @ts-ignore
|
||||
return cb(clientConn, ...args);
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
onCall(
|
||||
"invokeFunction",
|
||||
(plugName: string, name: string, ...args: any[]): Promise<any> => {
|
||||
let plug = this.system.loadedPlugs.get(plugName);
|
||||
if (!plug) {
|
||||
throw new Error(`Plug ${plugName} not loaded`);
|
||||
}
|
||||
console.log(
|
||||
"Invoking function",
|
||||
name,
|
||||
"for plug",
|
||||
plugName,
|
||||
"as requested over socket"
|
||||
);
|
||||
return plug.invoke(name, args);
|
||||
}
|
||||
);
|
||||
|
||||
console.log("Sending the sytem to the client");
|
||||
socket.emit("loadSystem", this.system.toJSON());
|
||||
});
|
||||
}
|
||||
|
||||
close() {
|
||||
console.log("Closing server");
|
||||
(this.apis.get("index")! as IndexApi).db.destroy().catch((err) => {
|
||||
console.error(err);
|
||||
});
|
||||
}
|
||||
}
|
||||
+48
-2
@@ -1,8 +1,54 @@
|
||||
import { mkdir, readdir, readFile, stat, unlink, writeFile } from "fs/promises";
|
||||
import * as path from "path";
|
||||
import { PageMeta } from "./types";
|
||||
import { EventHook } from "../plugos/hooks/event";
|
||||
|
||||
export class DiskStorage {
|
||||
export interface Storage {
|
||||
listPages(): Promise<PageMeta[]>;
|
||||
|
||||
readPage(pageName: string): Promise<{ text: string; meta: PageMeta }>;
|
||||
|
||||
writePage(pageName: string, text: string): Promise<PageMeta>;
|
||||
|
||||
getPageMeta(pageName: string): Promise<PageMeta>;
|
||||
|
||||
deletePage(pageName: string): Promise<void>;
|
||||
}
|
||||
|
||||
export class EventedStorage implements Storage {
|
||||
constructor(private wrapped: Storage, private eventHook: EventHook) {}
|
||||
|
||||
listPages(): Promise<PageMeta[]> {
|
||||
return this.wrapped.listPages();
|
||||
}
|
||||
|
||||
readPage(pageName: string): Promise<{ text: string; meta: PageMeta }> {
|
||||
return this.wrapped.readPage(pageName);
|
||||
}
|
||||
|
||||
async writePage(pageName: string, text: string): Promise<PageMeta> {
|
||||
const newPageMeta = this.wrapped.writePage(pageName, text);
|
||||
// This can happen async
|
||||
this.eventHook.dispatchEvent("page:saved", pageName).then(() => {
|
||||
return this.eventHook.dispatchEvent("page:index", {
|
||||
name: pageName,
|
||||
text: text,
|
||||
});
|
||||
});
|
||||
return newPageMeta;
|
||||
}
|
||||
|
||||
getPageMeta(pageName: string): Promise<PageMeta> {
|
||||
return this.wrapped.getPageMeta(pageName);
|
||||
}
|
||||
|
||||
async deletePage(pageName: string): Promise<void> {
|
||||
await this.eventHook.dispatchEvent("page:deleted", pageName);
|
||||
return this.wrapped.deletePage(pageName);
|
||||
}
|
||||
}
|
||||
|
||||
export class DiskStorage implements Storage {
|
||||
rootPath: string;
|
||||
|
||||
constructor(rootPath: string) {
|
||||
@@ -88,7 +134,7 @@ export class DiskStorage {
|
||||
}
|
||||
}
|
||||
|
||||
async deletePage(pageName: string) {
|
||||
async deletePage(pageName: string): Promise<void> {
|
||||
let localPath = path.join(this.rootPath, pageName + ".md");
|
||||
await unlink(localPath);
|
||||
}
|
||||
|
||||
+196
-5
@@ -1,13 +1,26 @@
|
||||
import { Express } from "express";
|
||||
import express, { Express } from "express";
|
||||
import { SilverBulletHooks } from "../common/manifest";
|
||||
import { EndpointHook } from "../plugos/hooks/endpoint";
|
||||
import { readFile } from "fs/promises";
|
||||
import { System } from "../plugos/system";
|
||||
import cors from "cors";
|
||||
import { DiskStorage, EventedStorage, Storage } from "./disk_storage";
|
||||
import path from "path";
|
||||
import bodyParser from "body-parser";
|
||||
import { EventHook } from "../plugos/hooks/event";
|
||||
import spaceSyscalls from "./syscalls/space";
|
||||
import { eventSyscalls } from "../plugos/syscalls/event";
|
||||
import { pageIndexSyscalls } from "./syscalls";
|
||||
import knex, { Knex } from "knex";
|
||||
|
||||
export class ExpressServer {
|
||||
app: Express;
|
||||
system: System<SilverBulletHooks>;
|
||||
private rootPath: string;
|
||||
private storage: Storage;
|
||||
private distDir: string;
|
||||
private eventHook: EventHook;
|
||||
private db: Knex<any, unknown[]>;
|
||||
|
||||
constructor(
|
||||
app: Express,
|
||||
@@ -17,19 +30,197 @@ export class ExpressServer {
|
||||
) {
|
||||
this.app = app;
|
||||
this.rootPath = rootPath;
|
||||
this.distDir = distDir;
|
||||
this.system = system;
|
||||
|
||||
// Setup system
|
||||
this.eventHook = new EventHook();
|
||||
system.addHook(this.eventHook);
|
||||
this.storage = new EventedStorage(
|
||||
new DiskStorage(rootPath),
|
||||
this.eventHook
|
||||
);
|
||||
this.db = knex({
|
||||
client: "better-sqlite3",
|
||||
connection: {
|
||||
filename: path.join(rootPath, "data.db"),
|
||||
},
|
||||
useNullAsDefault: true,
|
||||
});
|
||||
system.registerSyscalls("index", [], pageIndexSyscalls(this.db));
|
||||
system.registerSyscalls("space", [], spaceSyscalls(this.storage));
|
||||
system.registerSyscalls("event", [], eventSyscalls(this.eventHook));
|
||||
system.addHook(new EndpointHook(app, "/_"));
|
||||
}
|
||||
|
||||
async init() {
|
||||
console.log("Setting up router");
|
||||
|
||||
let fsRouter = express.Router();
|
||||
|
||||
// Page list
|
||||
fsRouter.route("/").get(async (req, res) => {
|
||||
res.json(await this.storage.listPages());
|
||||
});
|
||||
|
||||
fsRouter.route("/").post(bodyParser.json(), async (req, res) => {});
|
||||
|
||||
fsRouter
|
||||
.route(/\/(.+)/)
|
||||
.get(async (req, res) => {
|
||||
let pageName = req.params[0];
|
||||
console.log("Getting", pageName);
|
||||
try {
|
||||
let pageData = await this.storage.readPage(pageName);
|
||||
res.status(200);
|
||||
res.header("Last-Modified", "" + pageData.meta.lastModified);
|
||||
res.header("Content-Type", "text/markdown");
|
||||
res.send(pageData.text);
|
||||
} catch (e) {
|
||||
res.status(200);
|
||||
res.send("");
|
||||
}
|
||||
})
|
||||
.put(bodyParser.text({ type: "*/*" }), async (req, res) => {
|
||||
let pageName = req.params[0];
|
||||
console.log("Saving", pageName);
|
||||
|
||||
try {
|
||||
let meta = await this.storage.writePage(pageName, 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 pageName = req.params[0];
|
||||
try {
|
||||
const meta = await this.storage.getPageMeta(pageName);
|
||||
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 pageName = req.params[0];
|
||||
try {
|
||||
await this.storage.deletePage(pageName);
|
||||
res.status(200);
|
||||
res.send("OK");
|
||||
} catch (e) {
|
||||
console.error("Error deleting file", e);
|
||||
res.status(500);
|
||||
res.send("OK");
|
||||
}
|
||||
});
|
||||
|
||||
this.app.use(
|
||||
"/fs",
|
||||
cors({
|
||||
methods: "GET,HEAD,PUT,OPTIONS,POST,DELETE",
|
||||
preflightContinue: true,
|
||||
}),
|
||||
fsRouter
|
||||
);
|
||||
|
||||
let plugRouter = express.Router();
|
||||
|
||||
// Plug list
|
||||
plugRouter.get("/", async (req, res) => {
|
||||
res.json(
|
||||
[...this.system.loadedPlugs.values()].map(({ name, version }) => ({
|
||||
name,
|
||||
version,
|
||||
}))
|
||||
);
|
||||
});
|
||||
|
||||
plugRouter.get("/:name", async (req, res) => {
|
||||
const plugName = req.params.name;
|
||||
const plug = this.system.loadedPlugs.get(plugName);
|
||||
if (!plug) {
|
||||
res.status(404);
|
||||
res.send("Not found");
|
||||
} else {
|
||||
res.header("Last-Modified", "" + plug.version);
|
||||
res.send(plug.manifest);
|
||||
}
|
||||
});
|
||||
plugRouter.post(
|
||||
"/:plug/syscall/:name",
|
||||
bodyParser.json(),
|
||||
async (req, res) => {
|
||||
const name = req.params.name;
|
||||
const plugName = req.params.plug;
|
||||
const args = req.body as any;
|
||||
const plug = this.system.loadedPlugs.get(plugName);
|
||||
if (!plug) {
|
||||
res.status(404);
|
||||
return res.send(`Plug ${plugName} not found`);
|
||||
}
|
||||
try {
|
||||
const result = await this.system.syscallWithContext(
|
||||
{ plug },
|
||||
name,
|
||||
args
|
||||
);
|
||||
res.status(200);
|
||||
res.send(result);
|
||||
} catch (e: any) {
|
||||
res.status(500);
|
||||
return res.send(e.message);
|
||||
}
|
||||
}
|
||||
);
|
||||
plugRouter.post(
|
||||
"/:plug/function/:name",
|
||||
bodyParser.json(),
|
||||
async (req, res) => {
|
||||
const name = req.params.name;
|
||||
const plugName = req.params.plug;
|
||||
const args = req.body as any[];
|
||||
const plug = this.system.loadedPlugs.get(plugName);
|
||||
if (!plug) {
|
||||
res.status(404);
|
||||
return res.send(`Plug ${plugName} not found`);
|
||||
}
|
||||
try {
|
||||
console.log("Invoking", name, "with args", args);
|
||||
const result = await plug.invoke(name, args);
|
||||
res.status(200);
|
||||
res.send(result);
|
||||
} catch (e: any) {
|
||||
res.status(500);
|
||||
console.log("Error invoking function", e);
|
||||
return res.send(e.message);
|
||||
}
|
||||
}
|
||||
);
|
||||
|
||||
this.app.use(
|
||||
"/plug",
|
||||
cors({
|
||||
methods: "GET,HEAD,PUT,OPTIONS,POST,DELETE",
|
||||
preflightContinue: true,
|
||||
}),
|
||||
plugRouter
|
||||
);
|
||||
|
||||
// Fallback, serve index.html
|
||||
let cachedIndex: string | undefined = undefined;
|
||||
app.get("/*", async (req, res) => {
|
||||
this.app.get("/*", async (req, res) => {
|
||||
if (!cachedIndex) {
|
||||
cachedIndex = await readFile(`${distDir}/index.html`, "utf8");
|
||||
cachedIndex = await readFile(`${this.distDir}/index.html`, "utf8");
|
||||
}
|
||||
res.status(200).header("Content-Type", "text/html").send(cachedIndex);
|
||||
});
|
||||
}
|
||||
|
||||
async init() {}
|
||||
}
|
||||
|
||||
@@ -1,84 +0,0 @@
|
||||
import { ApiProvider, ClientConnection } from "./api_server";
|
||||
import knex, { Knex } from "knex";
|
||||
import path from "path";
|
||||
import { ensurePageIndexTable, pageIndexSyscalls } from "./syscalls/page_index";
|
||||
|
||||
type IndexItem = {
|
||||
page: string;
|
||||
key: string;
|
||||
value: any;
|
||||
};
|
||||
|
||||
export class IndexApi implements ApiProvider {
|
||||
db: Knex<any, unknown>;
|
||||
|
||||
constructor(rootPath: string) {
|
||||
this.db = knex({
|
||||
client: "better-sqlite3",
|
||||
connection: {
|
||||
filename: path.join(rootPath, "data.db"),
|
||||
},
|
||||
useNullAsDefault: true,
|
||||
});
|
||||
}
|
||||
|
||||
async init() {
|
||||
await ensurePageIndexTable(this.db);
|
||||
}
|
||||
|
||||
api() {
|
||||
const syscalls = pageIndexSyscalls(this.db);
|
||||
const nullContext = { plug: null };
|
||||
return {
|
||||
clearPageIndexForPage: async (
|
||||
clientConn: ClientConnection,
|
||||
page: string
|
||||
) => {
|
||||
console.log("Now going to clear index for", page);
|
||||
return syscalls.clearPageIndexForPage(nullContext, page);
|
||||
},
|
||||
set: async (
|
||||
clientConn: ClientConnection,
|
||||
page: string,
|
||||
key: string,
|
||||
value: any
|
||||
) => {
|
||||
return syscalls.set(nullContext, page, key, value);
|
||||
},
|
||||
get: async (clientConn: ClientConnection, page: string, key: string) => {
|
||||
return syscalls.get(nullContext, page, key);
|
||||
},
|
||||
delete: async (
|
||||
clientConn: ClientConnection,
|
||||
page: string,
|
||||
key: string
|
||||
) => {
|
||||
return syscalls.delete(nullContext, page, key);
|
||||
},
|
||||
scanPrefixForPage: async (
|
||||
clientConn: ClientConnection,
|
||||
page: string,
|
||||
prefix: string
|
||||
) => {
|
||||
return syscalls.scanPrefixForPage(nullContext, page, prefix);
|
||||
},
|
||||
scanPrefixGlobal: async (
|
||||
clientConn: ClientConnection,
|
||||
prefix: string
|
||||
) => {
|
||||
return syscalls.scanPrefixGlobal(nullContext, prefix);
|
||||
},
|
||||
deletePrefixForPage: async (
|
||||
clientConn: ClientConnection,
|
||||
page: string,
|
||||
prefix: string
|
||||
) => {
|
||||
return syscalls.deletePrefixForPage(nullContext, page, prefix);
|
||||
},
|
||||
|
||||
clearPageIndex: async (clientConn: ClientConnection) => {
|
||||
return syscalls.clearPageIndex(nullContext);
|
||||
},
|
||||
};
|
||||
}
|
||||
}
|
||||
@@ -1,344 +0,0 @@
|
||||
import { ClientPageState, Page, PageMeta } from "./types";
|
||||
import { ChangeSet } from "@codemirror/state";
|
||||
import { Update } from "@codemirror/collab";
|
||||
import { ApiProvider, ClientConnection } from "./api_server";
|
||||
import { Socket } from "socket.io";
|
||||
import { DiskStorage } from "./disk_storage";
|
||||
import { safeRun } from "./util";
|
||||
import fs from "fs";
|
||||
import path from "path";
|
||||
import { stat } from "fs/promises";
|
||||
import { Cursor, cursorEffect } from "../webapp/cursorEffect";
|
||||
import { SilverBulletHooks } from "../common/manifest";
|
||||
import { System } from "../plugos/system";
|
||||
import { EventHook } from "../plugos/hooks/event";
|
||||
import spaceSyscalls from "./syscalls/space";
|
||||
import { eventSyscalls } from "../plugos/syscalls/event";
|
||||
|
||||
export class PageApi implements ApiProvider {
|
||||
openPages: Map<string, Page>;
|
||||
pageStore: DiskStorage;
|
||||
rootPath: string;
|
||||
connectedSockets: Set<Socket>;
|
||||
private system: System<SilverBulletHooks>;
|
||||
private eventHook: EventHook;
|
||||
|
||||
constructor(
|
||||
rootPath: string,
|
||||
connectedSockets: Set<Socket>,
|
||||
openPages: Map<string, Page>,
|
||||
system: System<SilverBulletHooks>
|
||||
) {
|
||||
this.pageStore = new DiskStorage(rootPath);
|
||||
this.rootPath = rootPath;
|
||||
this.openPages = openPages;
|
||||
this.connectedSockets = connectedSockets;
|
||||
this.system = system;
|
||||
this.eventHook = new EventHook();
|
||||
system.addHook(this.eventHook);
|
||||
system.registerSyscalls("space", [], spaceSyscalls(this));
|
||||
system.registerSyscalls("event", [], eventSyscalls(this.eventHook));
|
||||
}
|
||||
|
||||
async init(): Promise<void> {
|
||||
this.fileWatcher();
|
||||
// TODO: Move this elsewhere, this doesn't belong here
|
||||
this.system.on({
|
||||
plugLoaded: (plugName, plugDef) => {
|
||||
console.log("Plug updated on disk, broadcasting to all clients");
|
||||
this.connectedSockets.forEach((socket) => {
|
||||
socket.emit("plugLoaded", plugName, plugDef.manifest);
|
||||
});
|
||||
},
|
||||
plugUnloaded: (plugName) => {
|
||||
console.log("Plug removed on disk, broadcasting to all clients");
|
||||
this.connectedSockets.forEach((socket) => {
|
||||
socket.emit("plugUnloaded", plugName);
|
||||
});
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
broadcastCursors(page: Page) {
|
||||
page.clientStates.forEach((client) => {
|
||||
client.socket.emit(
|
||||
"cursorSnapshot",
|
||||
page.name,
|
||||
Object.fromEntries(page.cursors.entries())
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
async flushPageToDisk(name: string, page: Page) {
|
||||
let meta = await this.pageStore.writePage(name, page.text.sliceString(0));
|
||||
console.log(`Wrote page ${name} to disk`);
|
||||
page.meta = meta;
|
||||
}
|
||||
|
||||
async disconnectClient(client: ClientPageState, page: Page) {
|
||||
console.log("Disconnecting client");
|
||||
page.clientStates.delete(client);
|
||||
if (page.clientStates.size === 0) {
|
||||
console.log("No more clients for", page.name, "flushing");
|
||||
await this.flushPageToDisk(page.name, page);
|
||||
this.openPages.delete(page.name);
|
||||
} else {
|
||||
page.cursors.delete(client.socket.id);
|
||||
this.broadcastCursors(page);
|
||||
}
|
||||
}
|
||||
|
||||
fileWatcher() {
|
||||
fs.watch(
|
||||
this.rootPath,
|
||||
{
|
||||
recursive: true,
|
||||
persistent: false,
|
||||
},
|
||||
(eventType, filename) => {
|
||||
safeRun(async () => {
|
||||
if (!filename.endsWith(".md")) {
|
||||
return;
|
||||
}
|
||||
let localPath = path.join(this.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",
|
||||
pageName
|
||||
);
|
||||
for (let socket of this.connectedSockets) {
|
||||
socket.emit("pageCreated", {
|
||||
name: pageName,
|
||||
lastModified: modifiedTime,
|
||||
} as PageMeta);
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
api() {
|
||||
return {
|
||||
openPage: async (clientConn: ClientConnection, pageName: string) => {
|
||||
let page = this.openPages.get(pageName);
|
||||
if (!page) {
|
||||
try {
|
||||
let { text, meta } = await this.pageStore.readPage(pageName);
|
||||
page = new Page(pageName, text, meta);
|
||||
} catch (e) {
|
||||
console.log("Creating new page", pageName);
|
||||
page = new Page(pageName, "", { name: pageName, lastModified: 0 });
|
||||
}
|
||||
this.openPages.set(pageName, page);
|
||||
}
|
||||
page.clientStates.add(
|
||||
new ClientPageState(clientConn.sock, page.version)
|
||||
);
|
||||
clientConn.openPages.add(pageName);
|
||||
console.log("Opened page", pageName);
|
||||
this.broadcastCursors(page);
|
||||
return page.toJSON();
|
||||
},
|
||||
pushUpdates: async (
|
||||
clientConn: ClientConnection,
|
||||
pageName: string,
|
||||
version: number,
|
||||
updates: any[]
|
||||
): Promise<boolean> => {
|
||||
let page = this.openPages.get(pageName);
|
||||
|
||||
if (!page) {
|
||||
console.error(
|
||||
"Received updates for not open page",
|
||||
pageName,
|
||||
this.openPages.keys()
|
||||
);
|
||||
return false;
|
||||
}
|
||||
if (version !== page.version) {
|
||||
console.error("Invalid version", version, page.version);
|
||||
return false;
|
||||
} else {
|
||||
console.log("Applying", updates.length, "updates to", pageName);
|
||||
let transformedUpdates = [];
|
||||
let textChanged = false;
|
||||
for (let update of updates) {
|
||||
let changes = ChangeSet.fromJSON(update.changes);
|
||||
let transformedUpdate = {
|
||||
changes,
|
||||
clientID: update.clientID,
|
||||
effects: update.cursors?.map((c: Cursor) => {
|
||||
page!.cursors.set(c.userId, c);
|
||||
return cursorEffect.of(c);
|
||||
}),
|
||||
};
|
||||
page.updates.push(transformedUpdate);
|
||||
transformedUpdates.push(transformedUpdate);
|
||||
let oldText = page.text;
|
||||
page.text = changes.apply(page.text);
|
||||
if (oldText !== page.text) {
|
||||
textChanged = true;
|
||||
}
|
||||
}
|
||||
console.log(
|
||||
"New version",
|
||||
page.version,
|
||||
"Updates buffered:",
|
||||
page.updates.length
|
||||
);
|
||||
|
||||
if (textChanged) {
|
||||
// Throttle
|
||||
if (!page.saveTimer) {
|
||||
page.saveTimer = setTimeout(() => {
|
||||
safeRun(async () => {
|
||||
if (page) {
|
||||
console.log(
|
||||
"Persisting",
|
||||
pageName,
|
||||
" to disk and indexing."
|
||||
);
|
||||
await this.flushPageToDisk(pageName, page);
|
||||
await this.eventHook.dispatchEvent("page:saved", pageName);
|
||||
await this.eventHook.dispatchEvent("page:index", {
|
||||
name: pageName,
|
||||
text: page.text.sliceString(0),
|
||||
});
|
||||
page.saveTimer = undefined;
|
||||
}
|
||||
});
|
||||
}, 1000);
|
||||
}
|
||||
}
|
||||
while (page.pending.length) {
|
||||
page.pending.pop()!(transformedUpdates);
|
||||
}
|
||||
return true;
|
||||
}
|
||||
},
|
||||
|
||||
pullUpdates: async (
|
||||
clientConn: ClientConnection,
|
||||
pageName: string,
|
||||
version: number
|
||||
): Promise<Update[]> => {
|
||||
let page = this.openPages.get(pageName);
|
||||
// console.log("Pulling updates for", pageName);
|
||||
if (!page) {
|
||||
console.error("Fetching updates for not open page");
|
||||
return [];
|
||||
}
|
||||
// TODO: Optimize this
|
||||
let oldestVersion = Infinity;
|
||||
page.clientStates.forEach((client) => {
|
||||
oldestVersion = Math.min(client.version, oldestVersion);
|
||||
if (client.socket === clientConn.sock) {
|
||||
client.version = version;
|
||||
}
|
||||
});
|
||||
page.flushUpdates(oldestVersion);
|
||||
if (version < page.version) {
|
||||
return page.updatesSince(version);
|
||||
} else {
|
||||
return new Promise((resolve) => {
|
||||
page!.pending.push(resolve);
|
||||
});
|
||||
}
|
||||
},
|
||||
|
||||
readPage: async (
|
||||
clientConn: ClientConnection,
|
||||
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);
|
||||
}
|
||||
},
|
||||
|
||||
writePage: async (
|
||||
clientConn: ClientConnection,
|
||||
pageName: string,
|
||||
text: string
|
||||
) => {
|
||||
// Write to disk
|
||||
let pageMeta = await this.pageStore.writePage(pageName, text);
|
||||
|
||||
// Notify clients that have the page open
|
||||
let page = this.openPages.get(pageName);
|
||||
if (page) {
|
||||
for (let client of page.clientStates) {
|
||||
client.socket.emit("pageChanged", pageMeta);
|
||||
}
|
||||
this.openPages.delete(pageName);
|
||||
}
|
||||
// Trigger system events
|
||||
await this.eventHook.dispatchEvent("page:saved", pageName);
|
||||
await this.eventHook.dispatchEvent("page:index", {
|
||||
name: pageName,
|
||||
text: text,
|
||||
});
|
||||
return pageMeta;
|
||||
},
|
||||
|
||||
deletePage: async (clientConn: ClientConnection, pageName: string) => {
|
||||
this.openPages.delete(pageName);
|
||||
clientConn.openPages.delete(pageName);
|
||||
// Cascading of this to all connected clients will be handled by file watcher
|
||||
await this.pageStore.deletePage(pageName);
|
||||
await this.eventHook.dispatchEvent("page:deleted", pageName);
|
||||
},
|
||||
|
||||
listPages: async (clientConn: ClientConnection): Promise<PageMeta[]> => {
|
||||
return this.pageStore.listPages();
|
||||
},
|
||||
|
||||
getPageMeta: async (
|
||||
clientConn: ClientConnection,
|
||||
pageName: string
|
||||
): Promise<PageMeta> => {
|
||||
let page = this.openPages.get(pageName);
|
||||
if (page) {
|
||||
return page.meta;
|
||||
}
|
||||
return this.pageStore.getPageMeta(pageName);
|
||||
},
|
||||
};
|
||||
}
|
||||
}
|
||||
@@ -2,8 +2,6 @@
|
||||
|
||||
import express from "express";
|
||||
import http from "http";
|
||||
import {Server} from "socket.io";
|
||||
import {SocketServer} from "./api_server";
|
||||
import yargs from "yargs";
|
||||
import {hideBin} from "yargs/helpers";
|
||||
import {SilverBulletHooks} from "../common/manifest";
|
||||
@@ -31,23 +29,11 @@ const app = express();
|
||||
const server = http.createServer(app);
|
||||
const system = new System<SilverBulletHooks>("server");
|
||||
|
||||
const io = new Server(server, {
|
||||
cors: {
|
||||
methods: "GET,HEAD,PUT,OPTIONS,POST,DELETE",
|
||||
preflightContinue: true,
|
||||
},
|
||||
});
|
||||
|
||||
const port = args.port;
|
||||
const distDir = `${__dirname}/../webapp`;
|
||||
|
||||
app.use("/", express.static(distDir));
|
||||
|
||||
let socketServer = new SocketServer(pagesPath, io, system);
|
||||
socketServer.init().catch((e) => {
|
||||
console.error(e);
|
||||
});
|
||||
|
||||
const expressServer = new ExpressServer(app, pagesPath, distDir, system);
|
||||
expressServer
|
||||
.init()
|
||||
|
||||
@@ -1,27 +1,23 @@
|
||||
import { PageMeta } from "../types";
|
||||
import { SysCallMapping } from "../../plugos/system";
|
||||
import { PageApi } from "../page_api";
|
||||
import { ClientConnection } from "../api_server";
|
||||
import { Storage } from "../disk_storage";
|
||||
|
||||
export default (pageApi: PageApi): SysCallMapping => {
|
||||
const api = pageApi.api();
|
||||
// @ts-ignore
|
||||
const dummyConn = new ClientConnection(null);
|
||||
export default (storage: Storage): SysCallMapping => {
|
||||
return {
|
||||
listPages: (ctx): Promise<PageMeta[]> => {
|
||||
return api.listPages(dummyConn);
|
||||
return storage.listPages();
|
||||
},
|
||||
readPage: async (
|
||||
ctx,
|
||||
name: string
|
||||
): Promise<{ text: string; meta: PageMeta }> => {
|
||||
return api.readPage(dummyConn, name);
|
||||
return storage.readPage(name);
|
||||
},
|
||||
writePage: async (ctx, name: string, text: string): Promise<PageMeta> => {
|
||||
return api.writePage(dummyConn, name, text);
|
||||
return storage.writePage(name, text);
|
||||
},
|
||||
deletePage: async (ctx, name: string) => {
|
||||
return api.deletePage(dummyConn, name);
|
||||
return storage.deletePage(name);
|
||||
},
|
||||
};
|
||||
};
|
||||
|
||||
@@ -1,62 +1,5 @@
|
||||
import { Update } from "@codemirror/collab";
|
||||
import { Text } from "@codemirror/state";
|
||||
import { Socket } from "socket.io";
|
||||
import { Cursor } from "../webapp/cursorEffect";
|
||||
export class ClientPageState {
|
||||
constructor(public socket: Socket, public version: number) {}
|
||||
}
|
||||
|
||||
export type PageMeta = {
|
||||
name: string;
|
||||
lastModified: number;
|
||||
version?: number;
|
||||
};
|
||||
|
||||
export class Page {
|
||||
versionOffset = 0;
|
||||
updates: Update[] = [];
|
||||
cursors = new Map<string, Cursor>();
|
||||
clientStates = new Set<ClientPageState>();
|
||||
|
||||
pending: ((value: any) => void)[] = [];
|
||||
|
||||
text: Text;
|
||||
meta: PageMeta;
|
||||
|
||||
saveTimer: NodeJS.Timeout | undefined;
|
||||
name: string;
|
||||
|
||||
constructor(name: string, text: string, meta: PageMeta) {
|
||||
this.name = name;
|
||||
this.text = Text.of(text.split("\n"));
|
||||
this.meta = meta;
|
||||
}
|
||||
|
||||
updatesSince(version: number): Update[] {
|
||||
return this.updates.slice(version - this.versionOffset);
|
||||
}
|
||||
|
||||
get version(): number {
|
||||
return this.updates.length + this.versionOffset;
|
||||
}
|
||||
|
||||
flushUpdates(version: number) {
|
||||
if (this.versionOffset > version) {
|
||||
throw Error("This should never happen");
|
||||
}
|
||||
if (this.versionOffset === version) {
|
||||
return;
|
||||
}
|
||||
this.updates = this.updates.slice(version - this.versionOffset);
|
||||
this.versionOffset = version;
|
||||
// console.log("Flushed updates, now got", this.updates.length, "updates");
|
||||
}
|
||||
|
||||
toJSON() {
|
||||
return {
|
||||
text: this.text,
|
||||
version: this.version,
|
||||
cursors: Object.fromEntries(this.cursors.entries()),
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user