Compare commits

...

54 Commits

Author SHA1 Message Date
sriram veeraghanta 3e895da0e8 fix: live server setup 2025-04-04 16:42:27 +05:30
Surya Prashanth 483961ea5a chore: refactor redis, hocus-pocus implementations 2025-04-04 14:54:04 +05:30
sriram veeraghanta ecefd47823 fix: startup script fixes 2025-04-04 01:18:22 +05:30
sriram veeraghanta 3359172acf fix: live server setup and start process 2025-04-04 01:10:34 +05:30
sriram veeraghanta 18a3771315 fix: removed unused dependencies 2025-04-03 20:50:29 +05:30
sriram veeraghanta 31c09d5a48 Merge branch 'preview' of github.com:makeplane/plane into fix/live-server-restructuring 2025-04-03 20:25:16 +05:30
Palanikannan M 1315acf952 fix: redis extension fixes 2025-04-03 17:56:33 +05:30
Palanikannan M 728f517cb1 fix: logger reference 2025-04-03 17:47:56 +05:30
Palanikannan M cb3b6fbd8d Merge branch 'preview' into fix/live-server-restructuring 2025-04-03 17:37:38 +05:30
Palanikannan M 70dbad648d chore: remove services 2025-04-03 17:37:03 +05:30
Palanikannan M ad9888ce45 fix: registering handlers made simplified 2025-04-02 15:31:05 +05:30
Palanikannan M 4eef115fcd Merge branch 'preview' into fix/live-server-restructuring 2025-04-02 13:19:29 +05:30
Palanikannan M cabd5f4275 fix: moving around pages 2025-04-02 13:10:47 +05:30
Palanikannan M 7cbdcd4c94 fix: added zod validation 2025-03-26 19:58:32 +05:30
Palanikannan M cef60b55e4 fix: error factory types and document controller error handling 2025-03-26 15:52:36 +05:30
Palanikannan M 07bf68c8ca fix: exit with error 1 2025-03-26 14:55:41 +05:30
Palanikannan M 3d07d0f678 fix: add missing packages 2025-03-26 14:52:11 +05:30
Palanikannan M 6c2fb4b287 merge decorators 2025-03-26 14:40:55 +05:30
Palanikannan M c1b8feaf6f chore: Merge preview 2025-03-26 14:01:50 +05:30
Palanikannan M 8a9cdc6133 fix: refactor decorators 2025-03-26 13:59:28 +05:30
Palanikannan M 6db75fe29c fix: added package dependency 2025-03-26 00:05:19 +05:30
Palanikannan M ea7ebe66b1 feat: express decorators for rest apis and websocket 2025-03-25 20:30:47 +05:30
Palanikannan M 93066ef5d5 Merge branch 'fix/live-server-restructuring' into fix/live-server-restructuring 2025-03-25 19:46:58 +05:30
Palanikannan M 2db9d35678 Merge branch 'preview' into devin/1734544044-refactor-live-server 2025-03-25 19:42:44 +05:30
Palanikannan M 37cd01e306 feat: express decorators for rest apis and websocket 2025-03-25 17:51:53 +05:30
Palanikannan M d3defc9785 fix: redis process improved a lot 2025-03-24 22:03:59 +05:30
Palanikannan M 3a0891e0ee fix: extensions index ts cleaned up 2025-03-22 03:06:12 +05:30
Palanikannan M 53efad3399 chore: restructuring project-page methods, transformers and handlers 2025-03-22 02:23:37 +05:30
Palanikannan M 304ef1a80c Merge branch 'preview' into fix/live-server-restructuring 2025-03-22 01:51:37 +05:30
Palanikannan M 16d41a3841 fix: not shutting down the app for any reason 2025-03-22 01:51:28 +05:30
Palanikannan M a9f4427b21 fix: seperated decorators into it's own package 2025-03-21 03:04:19 +05:30
Palanikannan M e4f31aea08 Merge branch 'preview' into devin/1734544044-refactor-live-server 2025-03-19 16:01:26 +05:30
Palanikannan M c2a3e47d3d fix: stop event prop on error 2025-03-19 16:01:09 +05:30
Palanikannan M cef4110eb0 fix: handlers 2025-03-19 14:01:19 +05:30
Palanikannan M 3672ee4ef1 fix: errors and imports 2025-03-18 19:25:31 +05:30
Palanikannan M c56097b8c0 fix: better error handling for redis client 2025-03-18 18:56:50 +05:30
Palanikannan M 0d57e0ab32 fix: file structure for error handling 2025-03-18 17:32:46 +05:30
Palanikannan M df35ccecc9 fix: better error handling 2025-03-18 16:59:04 +05:30
Palanikannan M 38d8d3ea9b fix: dividing server code 2025-03-18 01:36:14 +05:30
Palanikannan M 3710b182d3 fix: tsup hot reloading 2025-03-18 00:16:30 +05:30
Palanikannan M 6897575a62 fix: logger and added working global error handling with sentry 2025-03-17 23:24:03 +05:30
Palanikannan M 388151b70b Merge branch 'preview' into devin/1734544044-refactor-live-server 2025-03-17 15:52:49 +05:30
Palanikannan M 3d61604569 Merge branch 'preview' into devin/1734544044-refactor-live-server 2025-02-08 20:32:50 +05:30
Palanikannan M 1b29f65664 fix: removed .js imports 2024-12-23 18:16:53 +05:30
Palanikannan M a229508611 Merge branch 'preview' into devin/1734544044-refactor-live-server 2024-12-23 17:23:59 +05:30
Devin AI d5bd4ef63a chore: switch from esbuild to tsup for build configuration
Co-Authored-By: sriram@plane.so <sriram@plane.so>
2024-12-19 07:49:27 +00:00
Devin AI 146332fff3 chore: replace babel with esbuild for build configuration
Co-Authored-By: sriram@plane.so <sriram@plane.so>
2024-12-19 07:42:19 +00:00
Devin AI 5802858772 fix: resolve typescript errors in server and decorators
Co-Authored-By: sriram@plane.so <sriram@plane.so>
2024-12-18 18:05:01 +00:00
Devin AI b39ce9c18a fix: add eslint config and fix websocket router interface
Co-Authored-By: sriram@plane.so <sriram@plane.so>
2024-12-18 18:03:29 +00:00
Devin AI a7ab5ae680 fix: resolve lint errors in decorators and collaboration controller
Co-Authored-By: sriram@plane.so <sriram@plane.so>
2024-12-18 18:02:48 +00:00
Devin AI f2a08853e2 feat: add start.ts entry point for live server
Co-Authored-By: sriram@plane.so <sriram@plane.so>
2024-12-18 18:01:12 +00:00
Devin AI 23eeb45713 feat: implement collaboration controller with websocket support
Co-Authored-By: sriram@plane.so <sriram@plane.so>
2024-12-18 17:54:15 +00:00
Devin AI 6c83a0df09 feat: implement health check controller with decorator support
Co-Authored-By: sriram@plane.so <sriram@plane.so>
2024-12-18 17:52:40 +00:00
Devin AI dbee7488e1 feat: add decorator system and controller registration for live server
Co-Authored-By: sriram@plane.so <sriram@plane.so>
2024-12-18 17:50:52 +00:00
54 changed files with 5370 additions and 4150 deletions
-23
View File
@@ -1,23 +0,0 @@
{
"presets": [
[
"@babel/preset-env",
{
"modules": false
}
],
"@babel/preset-typescript"
],
"plugins": [
[
"module-resolver",
{
"root": ["./src"],
"alias": {
"@/core": "./src/core",
"@/plane-live": "./src/ce"
}
}
]
]
}
+15 -15
View File
@@ -3,13 +3,15 @@
"version": "0.25.3",
"license": "AGPL-3.0",
"description": "A realtime collaborative server powers Plane's rich text editor",
"main": "./src/server.ts",
"main": "./dist/start.js",
"module": "./dist/start.mjs",
"types": "./dist/start.d.ts",
"private": true,
"type": "module",
"scripts": {
"dev": "PORT=3100 concurrently \"babel src --out-dir dist --extensions '.ts,.js' --watch\" \"nodemon dist/server.js\"",
"build": "babel src --out-dir dist --extensions \".ts,.js\"",
"start": "node dist/server.js",
"dev": "tsup --watch --onSuccess 'node --env-file=.env dist/start.js'",
"build": "tsup",
"start": "node --env-file=.env dist/start.js",
"lint": "eslint src --ext .ts,.tsx",
"lint:errors": "eslint src --ext .ts,.tsx --quiet"
},
@@ -21,14 +23,17 @@
"@hocuspocus/extension-redis": "^2.15.0",
"@hocuspocus/server": "^2.15.0",
"@plane/constants": "*",
"@plane/decorators": "*",
"@plane/editor": "*",
"@plane/logger": "*",
"@plane/types": "*",
"@tiptap/core": "2.10.4",
"@tiptap/html": "2.11.0",
"axios": "^1.8.3",
"compression": "^1.7.4",
"cookie-parser": "^1.4.7",
"cors": "^2.8.5",
"dotenv": "^16.4.5",
"dotenv": "^16.4.7",
"express": "^4.21.2",
"express-ws": "^5.0.2",
"helmet": "^7.1.0",
@@ -37,27 +42,22 @@
"morgan": "^1.10.0",
"pino-http": "^10.3.0",
"pino-pretty": "^11.2.2",
"reflect-metadata": "^0.2.2",
"uuid": "^10.0.0",
"y-prosemirror": "^1.2.15",
"y-protocols": "^1.0.6",
"yjs": "^13.6.20"
"yjs": "^13.6.20",
"zod": "^3.24.2"
},
"devDependencies": {
"@babel/cli": "^7.25.6",
"@babel/core": "^7.25.2",
"@babel/preset-env": "^7.25.4",
"@babel/preset-typescript": "^7.24.7",
"@types/compression": "^1.7.5",
"@types/cookie-parser": "^1.4.8",
"@types/cors": "^2.8.17",
"@types/dotenv": "^8.2.0",
"@types/express": "^4.17.21",
"@types/express-ws": "^3.0.4",
"@types/node": "^20.14.9",
"babel-plugin-module-resolver": "^5.0.2",
"concurrently": "^9.0.1",
"nodemon": "^3.1.7",
"ts-node": "^10.9.2",
"tsup": "^8.4.0",
"tsup": "8.4.0",
"typescript": "5.3.3"
}
}
+1
View File
@@ -0,0 +1 @@
export * from "./register";
@@ -0,0 +1,2 @@
export * from "./transformers";
export * from "./project-page-handler";
@@ -0,0 +1,74 @@
import {
DocumentHandler,
DocumentFetchParams,
DocumentStoreParams,
HandlerDefinition,
} from "@/core/types/document-handler";
import { handlerFactory } from "@/core/handlers/document-handlers/handler-factory";
import { PageService } from "@/services/page.service";
import { transformHTMLToBinary } from "./transformers";
import { getAllDocumentFormatsFromBinaryData } from "@/core/helpers/page";
const pageService = new PageService();
/**
* Handler for "project_page" document type
*/
export const projectPageHandler: DocumentHandler = {
/**
* Fetch project page description
*/
fetch: async ({ pageId, params, context }: DocumentFetchParams) => {
const { cookie } = context;
const workspaceSlug = params.get("workspaceSlug")?.toString();
const projectId = params.get("projectId")?.toString();
if (!workspaceSlug || !projectId || !cookie) return null;
const response = await pageService.fetchDescriptionBinary(workspaceSlug, projectId, pageId, cookie);
const binaryData = new Uint8Array(response);
if (binaryData.byteLength === 0) {
const binary = await transformHTMLToBinary(workspaceSlug, projectId, pageId, cookie);
if (binary) {
return binary;
}
}
return binaryData;
},
/**
* Store project page description
*/
store: async ({ pageId, state, params, context }: DocumentStoreParams) => {
const { cookie } = context;
if (!(state instanceof Uint8Array)) {
throw new Error("Invalid state: must be an instance of Uint8Array");
}
const workspaceSlug = params?.get("workspaceSlug")?.toString();
const projectId = params?.get("projectId")?.toString();
if (!workspaceSlug || !projectId || !cookie) return;
const { contentBinaryEncoded, contentHTML, contentJSON } = getAllDocumentFormatsFromBinaryData(state);
const payload = {
description_binary: contentBinaryEncoded,
description_html: contentHTML,
description: contentJSON,
};
await pageService.updateDescription(workspaceSlug, projectId, pageId, payload, cookie);
},
};
// Define the project page handler definition
export const projectPageHandlerDefinition: HandlerDefinition = {
selector: (context) => context.documentType === "project_page",
handler: projectPageHandler,
priority: 10, // Standard priority
};
// Register the handler directly from CE
export function registerProjectPageHandler() {
handlerFactory.register(projectPageHandlerDefinition);
}
@@ -0,0 +1,26 @@
import { PageService } from "@/services/page.service";
import { getBinaryDataFromHTMLString } from "@/core/helpers/page";
import { logger } from "@plane/logger";
const pageService = new PageService();
/**
* Transforms HTML description to binary format
*/
export const transformHTMLToBinary = async (
workspaceSlug: string,
projectId: string,
pageId: string,
cookie: string
) => {
if (!workspaceSlug || !projectId || !cookie) return;
try {
const pageDetails = await pageService.fetchDetails(workspaceSlug, projectId, pageId, cookie);
const { contentBinary } = getBinaryDataFromHTMLString(pageDetails.description_html ?? "<p></p>");
return contentBinary;
} catch (error) {
logger.error("Error while transforming from HTML to Uint8Array", error);
throw error;
}
};
+5
View File
@@ -0,0 +1,5 @@
import { registerProjectPageHandler } from "./project-page/project-page-handler";
export function initializeDocumentHandlers() {
registerProjectPageHandler();
}
-14
View File
@@ -1,14 +0,0 @@
// types
import { TDocumentTypes } from "@/core/types/common.js";
type TArgs = {
cookie: string | undefined;
documentType: TDocumentTypes | undefined;
pageId: string;
params: URLSearchParams;
}
export const fetchDocument = async (args: TArgs): Promise<Uint8Array | null> => {
const { documentType } = args;
throw Error(`Fetch failed: Invalid document type ${documentType} provided.`);
}
-15
View File
@@ -1,15 +0,0 @@
// types
import { TDocumentTypes } from "@/core/types/common.js";
type TArgs = {
cookie: string | undefined;
documentType: TDocumentTypes | undefined;
pageId: string;
params: URLSearchParams;
updatedDescription: Uint8Array;
}
export const updateDocument = async (args: TArgs): Promise<void> => {
const { documentType } = args;
throw Error(`Update failed: Invalid document type ${documentType} provided.`);
}
+1 -1
View File
@@ -1 +1 @@
export type TAdditionalDocumentTypes = {};
export type TAdditionalDocumentTypes = string;
@@ -0,0 +1,96 @@
import type { Request } from "express";
import type { WebSocket as WS } from "ws";
import type { Hocuspocus } from "@hocuspocus/server";
import { ErrorCategory } from "@/lib/error-handling/error-handler";
import { logger } from "@plane/logger";
import Errors from "@/lib/error-handling/error-factory";
import { Controller, WebSocket } from "@plane/decorators";
@Controller("/collaboration")
export class CollaborationController {
private metrics = {
errors: 0,
};
constructor(private readonly hocusPocusServer: Hocuspocus) {}
@WebSocket("/")
handleConnection(ws: WS, req: Request) {
const clientInfo = {
ip: req.ip,
userAgent: req.get("user-agent"),
requestId: req.id || crypto.randomUUID(),
};
try {
// Initialize the connection with Hocuspocus
this.hocusPocusServer.handleConnection(ws, req);
// Set up error handling for the connection
ws.on("error", (error) => {
this.handleConnectionError(error, clientInfo, ws);
});
} catch (error) {
this.handleConnectionError(error, clientInfo, ws);
}
}
private handleConnectionError(error: unknown, clientInfo: Record<string, any>, ws: WS) {
// Convert to AppError if needed
const appError = Errors.convertError(error instanceof Error ? error : new Error(String(error)), {
context: {
...clientInfo,
component: "WebSocketConnection",
},
});
// Log at appropriate level based on error category
if (appError.category === ErrorCategory.OPERATIONAL) {
logger.info(`WebSocket operational error: ${appError.message}`, {
error: appError,
clientInfo,
});
} else {
logger.error(`WebSocket error: ${appError.message}`, {
error: appError,
clientInfo,
stack: appError.stack,
});
}
// Alert if error threshold is reached
if (this.metrics.errors % 10 === 0) {
logger.warn(`High WebSocket error rate detected: ${this.metrics.errors} total errors`);
}
// Try to send error to client before closing
try {
if (ws.readyState === ws.OPEN) {
ws.send(
JSON.stringify({
type: "error",
message: appError.category === ErrorCategory.OPERATIONAL ? appError.message : "Internal server error",
})
);
}
} catch (sendError) {
// Ignore send errors at this point
}
// Close with informative message if connection is still open
if (ws.readyState === ws.OPEN) {
ws.close(
1011,
appError.category === ErrorCategory.OPERATIONAL
? `Error: ${appError.message}. Reconnect with exponential backoff.`
: "Internal server error. Please retry in a few moments."
);
}
}
getErrorMetrics() {
return {
errors: this.metrics.errors,
};
}
}
+128
View File
@@ -0,0 +1,128 @@
import type { Request, Response } from "express";
import { z } from "zod";
// helpers
import { convertHTMLDocumentToAllFormats } from "@/core/helpers/convert-document";
// types
import { TConvertDocumentRequestBody } from "@/core/types/common";
// decorators
import { CatchErrors } from "@/lib/error";
// logger
import { logger } from "@plane/logger";
import { Controller, Post } from "@plane/decorators";
import { AppError } from "@/lib/error-handling/error-handler";
import { handleError } from "@/lib/error-handling/error-factory";
// Define the schema with more robust validation
const convertDocumentSchema = z.object({
description_html: z
.string()
.min(1, "HTML content cannot be empty")
.refine((html) => html.trim().length > 0, "HTML content cannot be just whitespace")
.refine((html) => html.includes("<") && html.includes(">"), "Content must be valid HTML"),
variant: z.enum(["rich", "document"]),
});
@Controller("/convert-document")
export class DocumentController {
private metrics = {
conversions: 0,
errors: 0,
};
@Post("/")
@CatchErrors()
async convertDocument(req: Request, res: Response) {
const requestId = req.id || crypto.randomUUID();
const clientInfo = {
ip: req.ip,
userAgent: req.get("user-agent"),
requestId,
};
try {
// Validate request body
const validatedData = convertDocumentSchema.parse(req.body as TConvertDocumentRequestBody);
const { description_html, variant } = validatedData;
// Log validated data
logger.info("Validated document conversion request", {
...clientInfo,
variant,
contentLength: description_html.length,
});
// Process document conversion
const { description, description_binary } = convertHTMLDocumentToAllFormats({
document_html: description_html,
variant,
});
// Update metrics
this.metrics.conversions++;
// Log successful conversion
logger.info("Document conversion successful", {
...clientInfo,
variant,
outputLength: description_html.length,
});
// Return successful response
res.status(200).json({
description,
description_binary,
});
} catch (error) {
// Update error metrics
this.metrics.errors++;
let appError: AppError;
if (error instanceof z.ZodError) {
// Handle validation errors
appError = handleError(error, {
errorType: "unprocessable-entity",
message: "Invalid request data",
component: "document-conversion-controller",
operation: "convertDocument",
extraContext: {
...clientInfo,
validationErrors: error.errors.map((err) => ({
path: err.path.join("."),
message: err.message,
})),
},
});
} else {
// Handle other errors
appError = handleError(error, {
errorType: "internal",
message: "Internal server error",
component: "document-conversion-controller",
operation: "convertDocument",
extraContext: clientInfo,
});
}
// Log the error
logger.error("Document conversion failed", {
error: appError,
status: appError.status,
context: appError.context,
});
res.status(appError.status).json({
message: appError.message,
status: appError.status,
context: appError.context,
});
}
}
getMetrics() {
return {
conversions: this.metrics.conversions,
errors: this.metrics.errors,
};
}
}
+16
View File
@@ -0,0 +1,16 @@
import { CatchErrors } from "@/lib/error";
import { Controller, Get } from "@plane/decorators";
import type { Request, Response } from "express";
@Controller("/health")
export class HealthController {
@Get("/")
@CatchErrors()
async healthCheck(_req: Request, res: Response) {
res.status(200).json({
status: "OK",
timestamp: new Date().toISOString(),
version: process.env.APP_VERSION || "1.0.0",
});
}
}
+7
View File
@@ -0,0 +1,7 @@
import { HealthController } from "./health.controller";
import { DocumentController } from "./document.controller";
import { CollaborationController } from "./collaboration.controller";
export const REST_CONTROLLERS = [HealthController, DocumentController];
export const WEBSOCKET_CONTROLLERS = [CollaborationController];
+116
View File
@@ -0,0 +1,116 @@
import { Database } from "@hocuspocus/extension-database";
import { catchAsync } from "@/lib/error-handling/error-handler";
import { handleError } from "@/lib/error-handling/error-factory";
import { getDocumentHandler } from "../handlers/document-handlers";
import { type HocusPocusServerContext, type TDocumentTypes } from "@/core/types/common";
export const createDatabaseExtension = () => {
return new Database({
fetch: handleFetch,
store: handleStore,
});
};
const handleFetch = async ({
context,
documentName: pageId,
requestParameters,
}: {
context: HocusPocusServerContext;
documentName: TDocumentTypes;
requestParameters: URLSearchParams;
}) => {
const { documentType } = context;
const params = requestParameters;
let fetchedData = null;
fetchedData = await catchAsync(
async () => {
if (!documentType) {
handleError(null, {
errorType: "bad-request",
message: "Document type is required",
component: "database-extension",
operation: "fetch",
extraContext: { pageId },
throw: true,
});
}
const documentHandler = getDocumentHandler(documentType, context);
fetchedData = await documentHandler.fetch({
context: context as HocusPocusServerContext,
pageId,
params,
});
if (!fetchedData) {
handleError(null, {
errorType: "not-found",
message: `Failed to fetch document: ${pageId}`,
component: "database-extension",
operation: "fetch",
extraContext: { documentType, pageId },
});
}
return fetchedData;
},
{
params: { pageId, documentType: context.documentType },
extra: { operation: "fetch" },
}
)();
return fetchedData;
};
const handleStore = async ({
context,
state,
documentName: pageId,
requestParameters,
}: {
context: HocusPocusServerContext;
state: Buffer;
documentName: TDocumentTypes;
requestParameters: URLSearchParams;
}) => {
catchAsync(
async () => {
if (!state) {
handleError(null, {
errorType: "bad-request",
message: "Loaded binary state is required",
component: "database-extension",
operation: "store",
extraContext: { pageId },
throw: true,
});
}
const { documentType } = context as HocusPocusServerContext;
const params = requestParameters;
if (!documentType) {
handleError(null, {
errorType: "bad-request",
message: "Document type is required",
component: "database-extension",
operation: "store",
extraContext: { pageId },
throw: true,
});
}
const documentHandler = getDocumentHandler(documentType, context);
await documentHandler.store({
context: context as HocusPocusServerContext,
pageId,
state,
params,
});
},
{
params: { pageId, documentType: context.documentType },
extra: { operation: "store" },
}
)();
};
+12 -128
View File
@@ -1,142 +1,26 @@
// Third-party libraries
import { Redis } from "ioredis";
// Hocuspocus extensions and core
import { Database } from "@hocuspocus/extension-database";
// hocuspocus extensions and core
import { Extension } from "@hocuspocus/server";
import { Logger } from "@hocuspocus/extension-logger";
import { Redis as HocusPocusRedis } from "@hocuspocus/extension-redis";
// core helpers and utilities
import { manualLogger } from "@/core/helpers/logger.js";
import { getRedisUrl } from "@/core/lib/utils/redis-url.js";
// core libraries
import {
fetchPageDescriptionBinary,
updatePageDescription,
} from "@/core/lib/page.js";
// plane live libraries
import { fetchDocument } from "@/plane-live/lib/fetch-document.js";
import { updateDocument } from "@/plane-live/lib/update-document.js";
// types
import {
type HocusPocusServerContext,
type TDocumentTypes,
} from "@/core/types/common.js";
import { setupRedisExtension } from "@/core/extensions/redis";
import { createDatabaseExtension } from "@/core/extensions/database";
import { logger } from "@plane/logger";
export const getExtensions: () => Promise<Extension[]> = async () => {
export const getExtensions = async (): Promise<Extension[]> => {
const extensions: Extension[] = [
new Logger({
onChange: false,
log: (message) => {
manualLogger.info(message);
},
}),
new Database({
fetch: async ({ context, documentName: pageId, requestParameters }) => {
const cookie = (context as HocusPocusServerContext).cookie;
// query params
const params = requestParameters;
const documentType = params.get("documentType")?.toString() as
| TDocumentTypes
| undefined;
// TODO: Fix this lint error.
// eslint-disable-next-line no-async-promise-executor
return new Promise(async (resolve) => {
try {
let fetchedData = null;
if (documentType === "project_page") {
fetchedData = await fetchPageDescriptionBinary(
params,
pageId,
cookie,
);
} else {
fetchedData = await fetchDocument({
cookie,
documentType,
pageId,
params,
});
}
resolve(fetchedData);
} catch (error) {
manualLogger.error("Error in fetching document", error);
}
});
},
store: async ({
context,
state,
documentName: pageId,
requestParameters,
}) => {
const cookie = (context as HocusPocusServerContext).cookie;
// query params
const params = requestParameters;
const documentType = params.get("documentType")?.toString() as
| TDocumentTypes
| undefined;
// TODO: Fix this lint error.
// eslint-disable-next-line no-async-promise-executor
return new Promise(async () => {
try {
if (documentType === "project_page") {
await updatePageDescription(params, pageId, state, cookie);
} else {
await updateDocument({
cookie,
documentType,
pageId,
params,
updatedDescription: state,
});
}
} catch (error) {
manualLogger.error("Error in updating document:", error);
}
});
logger.info(message);
},
}),
createDatabaseExtension(),
];
const redisUrl = getRedisUrl();
if (redisUrl) {
try {
const redisClient = new Redis(redisUrl);
await new Promise<void>((resolve, reject) => {
redisClient.on("error", (error: any) => {
if (
error?.code === "ENOTFOUND" ||
error.message.includes("WRONGPASS") ||
error.message.includes("NOAUTH")
) {
redisClient.disconnect();
}
manualLogger.warn(
`Redis Client wasn't able to connect, continuing without Redis (you won't be able to sync data between multiple plane live servers)`,
error,
);
reject(error);
});
redisClient.on("ready", () => {
extensions.push(new HocusPocusRedis({ redis: redisClient }));
manualLogger.info("Redis Client connected ✅");
resolve();
});
});
} catch (error) {
manualLogger.warn(
`Redis Client wasn't able to connect, continuing without Redis (you won't be able to sync data between multiple plane live servers)`,
error,
);
}
} else {
manualLogger.warn(
"Redis URL is not set, continuing without Redis (you won't be able to sync data between multiple plane live servers)",
);
// Add Redis extensions if Redis is available
const redisExtension = await setupRedisExtension();
if (redisExtension) {
logger.info("HocusPocus Redis extension configured ✅");
extensions.push(redisExtension);
}
return extensions;
+16
View File
@@ -0,0 +1,16 @@
import { Redis as HocusPocusRedis } from "@hocuspocus/extension-redis";
// core helpers and utilities
import { logger } from "@plane/logger";
import { getRedisClient } from "@/core/lib/redis-manager";
/**
* Sets up the Redis extension for HocusPocus using the RedisManager singleton
* @returns Promise that resolves to a Redis extension array
*/
export const setupRedisExtension = () => {
// Wait for Redis connection
return new HocusPocusRedis({
redis: getRedisClient(),
});
};
@@ -0,0 +1,32 @@
import { DocumentHandler, HandlerContext, HandlerDefinition } from "@/core/types/document-handler";
/**
* Class that manages handler selection based on multiple criteria
*/
export class DocumentHandlerFactory {
private handlers: HandlerDefinition[] = [];
/**
* Register a handler with its selection criteria
*/
register(definition: HandlerDefinition): void {
this.handlers.push(definition);
// Sort handlers by priority (highest first)
this.handlers.sort((a, b) => b.priority - a.priority);
}
/**
* Get the appropriate handler based on the provided context
*/
getHandler(context: HandlerContext): DocumentHandler {
// Find the first handler whose selector returns true
const matchingHandler = this.handlers.find(h => h.selector(context));
// Return the matching handler or fall back to null/undefined
// (This will cause an error if no handlers match, which is good for debugging)
return matchingHandler?.handler as DocumentHandler;
}
}
// Create the singleton instance
export const handlerFactory = new DocumentHandlerFactory();
@@ -0,0 +1,31 @@
import { DocumentHandler, HandlerContext } from "@/core/types/document-handler";
import { handlerFactory } from "@/core/handlers/document-handlers/handler-factory";
import { HocusPocusServerContext } from "@/core/types/common";
import { initializeDocumentHandlers } from "@/plane-live/document-types";
// Initialize all CE document handlers
initializeDocumentHandlers();
/**
* Get a document handler based on the provided context criteria
* @param documentType The primary document type
* @param additionalContext Optional additional context criteria
* @returns The appropriate document handler
*/
export function getDocumentHandler(
documentType: string,
additionalContext: Omit<HocusPocusServerContext, "documentType">
): DocumentHandler {
// Create a context object with all criteria
const context: HandlerContext = {
documentType: documentType as any,
...additionalContext,
};
// Use the factory to get the appropriate handler
return handlerFactory.getHandler(context);
}
// Export the factory for direct access if needed
export { handlerFactory };
-21
View File
@@ -1,21 +0,0 @@
import { ErrorRequestHandler } from "express";
import { manualLogger } from "@/core/helpers/logger.js";
export const errorHandler: ErrorRequestHandler = (err, _req, res) => {
// Log the error
manualLogger.error(err);
// Set the response status
res.status(err.status || 500);
// Send the response
res.json({
error: {
message:
process.env.NODE_ENV === "production"
? "An unexpected error occurred"
: err.message,
...(process.env.NODE_ENV !== "production" && { stack: err.stack }),
},
});
};
+267
View File
@@ -0,0 +1,267 @@
import { handleError } from "../../lib/error-handling/error-factory";
/**
* A simple validation utility that integrates with our error system.
*
* This provides a fluent interface for validating data and throwing
* appropriate errors if validation fails.
*/
export class Validator<T> {
constructor(
private readonly data: T,
private readonly name: string = "data"
) {}
/**
* Ensures a value is defined (not undefined or null)
*/
required(message?: string): Validator<T> {
if (this.data === undefined || this.data === null) {
throw handleError(new ValidationError(this.name, message || `${this.name} is required`), {
errorType: "bad-request",
component: "validation",
operation: "validateRequired",
extraContext: { field: this.name },
throw: true,
});
}
return this;
}
/**
* Ensures a value is a string
*/
string(message?: string): Validator<T> {
if (typeof this.data !== "string") {
throw handleError(new ValidationError(this.name, message || `${this.name} must be a string`), {
errorType: "bad-request",
component: "validation",
operation: "validateString",
extraContext: { field: this.name },
throw: true,
});
}
return this;
}
/**
* Ensures a string is not empty
*/
notEmpty(message?: string): Validator<T> {
if (typeof this.data === "string" && this.data.trim() === "") {
throw handleError(new ValidationError(this.name, message || `${this.name} cannot be empty`), {
errorType: "bad-request",
component: "validation",
operation: "validateNonEmptyString",
extraContext: { field: this.name },
throw: true,
});
}
return this;
}
/**
* Ensures a value is a number
*/
number(message?: string): Validator<T> {
if (typeof this.data !== "number" || isNaN(this.data)) {
throw handleError(new ValidationError(this.name, message || `${this.name} must be a valid number`), {
errorType: "bad-request",
component: "validation",
operation: "validateNumber",
extraContext: { field: this.name },
throw: true,
});
}
return this;
}
/**
* Ensures an array is not empty
*/
nonEmptyArray(message?: string): Validator<T> {
if (!Array.isArray(this.data) || this.data.length === 0) {
throw handleError(new ValidationError(this.name, message || `${this.name} must be a non-empty array`), {
errorType: "bad-request",
component: "validation",
operation: "validateArray",
extraContext: { field: this.name },
throw: true,
});
}
return this;
}
/**
* Ensures a value matches a regular expression
*/
match(regex: RegExp, message?: string): Validator<T> {
if (typeof this.data !== "string" || !regex.test(this.data)) {
throw handleError(new ValidationError(this.name, message || `${this.name} has an invalid format`), {
errorType: "bad-request",
component: "validation",
operation: "validateFormat",
extraContext: { field: this.name, format: regex.toString() },
throw: true,
});
}
return this;
}
/**
* Ensures a value is one of the allowed values
*/
oneOf(allowedValues: any[], message?: string): Validator<T> {
if (!allowedValues.includes(this.data)) {
throw handleError(
new ValidationError(this.name, message || `${this.name} must be one of: ${allowedValues.join(", ")}`),
{
errorType: "bad-request",
component: "validation",
operation: "validateEnum",
extraContext: { field: this.name, allowedValues },
throw: true,
}
);
}
return this;
}
/**
* Custom validation function
*/
custom(validationFn: (value: T) => boolean, message?: string): Validator<T> {
if (!validationFn(this.data)) {
throw handleError(new ValidationError(this.name, message || `${this.name} is invalid`), {
errorType: "bad-request",
component: "validation",
operation: "validateCustom",
extraContext: { field: this.name },
throw: true,
});
}
return this;
}
/**
* Get the validated data
*/
get(): T {
return this.data;
}
}
/**
* Create a new validator for a value
*/
export const validate = <T>(data: T, name?: string): Validator<T> => {
return new Validator(data, name);
};
export default validate;
export class ValidationError extends Error {
constructor(
public name: string,
message: string
) {
super(message);
this.name = name;
}
}
export const validateRequired = (value: any, name: string, message?: string) => {
if (value === undefined || value === null) {
throw handleError(new ValidationError(name, message || `${name} is required`), {
errorType: "bad-request",
component: "validation",
operation: "validateRequired",
extraContext: { field: name },
throw: true,
});
}
};
export const validateString = (value: any, name: string, message?: string) => {
if (typeof value !== "string") {
throw handleError(new ValidationError(name, message || `${name} must be a string`), {
errorType: "bad-request",
component: "validation",
operation: "validateString",
extraContext: { field: name },
throw: true,
});
}
};
export const validateNonEmptyString = (value: string, name: string, message?: string) => {
if (!value.trim()) {
throw handleError(new ValidationError(name, message || `${name} cannot be empty`), {
errorType: "bad-request",
component: "validation",
operation: "validateNonEmptyString",
extraContext: { field: name },
throw: true,
});
}
};
export const validateNumber = (value: any, name: string, message?: string) => {
if (typeof value !== "number" || isNaN(value)) {
throw handleError(new ValidationError(name, message || `${name} must be a valid number`), {
errorType: "bad-request",
component: "validation",
operation: "validateNumber",
extraContext: { field: name },
throw: true,
});
}
};
export const validateArray = (value: any, name: string, message?: string) => {
if (!Array.isArray(value) || value.length === 0) {
throw handleError(new ValidationError(name, message || `${name} must be a non-empty array`), {
errorType: "bad-request",
component: "validation",
operation: "validateArray",
extraContext: { field: name },
throw: true,
});
}
};
export const validateFormat = (value: string, name: string, format: RegExp, message?: string) => {
if (!format.test(value)) {
throw handleError(new ValidationError(name, message || `${name} has an invalid format`), {
errorType: "bad-request",
component: "validation",
operation: "validateFormat",
extraContext: { field: name, format: format.toString() },
throw: true,
});
}
};
export const validateEnum = (value: any, name: string, allowedValues: any[], message?: string) => {
if (!allowedValues.includes(value)) {
throw handleError(new ValidationError(name, message || `${name} must be one of: ${allowedValues.join(", ")}`), {
errorType: "bad-request",
component: "validation",
operation: "validateEnum",
extraContext: { field: name, allowedValues },
throw: true,
});
}
};
export const validateCustom = (value: any, name: string, validator: (value: any) => boolean, message?: string) => {
if (!validator(value)) {
throw handleError(new ValidationError(name, message || `${name} is invalid`), {
errorType: "bad-request",
component: "validation",
operation: "validateCustom",
extraContext: { field: name },
throw: true,
});
}
};
-73
View File
@@ -1,73 +0,0 @@
import { Server } from "@hocuspocus/server";
import { v4 as uuidv4 } from "uuid";
// lib
import { handleAuthentication } from "@/core/lib/authentication.js";
// extensions
import { getExtensions } from "@/core/extensions/index.js";
import {
DocumentCollaborativeEvents,
TDocumentEventsServer,
} from "@plane/editor/lib";
// editor types
import { TUserDetails } from "@plane/editor";
// types
import { type HocusPocusServerContext } from "@/core/types/common.js";
export const getHocusPocusServer = async () => {
const extensions = await getExtensions();
const serverName = process.env.HOSTNAME || uuidv4();
return Server.configure({
name: serverName,
onAuthenticate: async ({
requestHeaders,
context,
// user id used as token for authentication
token,
}) => {
let cookie: string | undefined = undefined;
let userId: string | undefined = undefined;
// Extract cookie (fallback to request headers) and userId from token (for scenarios where
// the cookies are not passed in the request headers)
try {
const parsedToken = JSON.parse(token) as TUserDetails;
userId = parsedToken.id;
cookie = parsedToken.cookie;
} catch (error) {
// If token parsing fails, fallback to request headers
console.error("Token parsing failed, using request headers:", error);
} finally {
// If cookie is still not found, fallback to request headers
if (!cookie) {
cookie = requestHeaders.cookie?.toString();
}
}
if (!cookie || !userId) {
throw new Error("Credentials not provided");
}
// set cookie in context, so it can be used throughout the ws connection
(context as HocusPocusServerContext).cookie = cookie;
try {
await handleAuthentication({
cookie,
userId,
});
} catch (error) {
throw Error("Authentication unsuccessful!");
}
},
async onStateless({ payload, document }) {
// broadcast the client event (derived from the server event) to all the clients so that they can update their state
const response =
DocumentCollaborativeEvents[payload as TDocumentEventsServer].client;
if (response) {
document.broadcastStateless(response);
}
},
extensions,
debounce: 10000,
});
};
+27 -7
View File
@@ -1,27 +1,47 @@
// services
import { UserService } from "@/core/services/user.service.js";
// core helpers
import { manualLogger } from "@/core/helpers/logger.js";
import { UserService } from "@/services/user.service";
import { handleError } from "@/lib/error-handling/error-factory";
const userService = new UserService();
type Props = {
cookie: string;
userId: string;
workspaceSlug: string;
};
export const handleAuthentication = async (props: Props) => {
const { cookie, userId } = props;
const { cookie, userId, workspaceSlug } = props;
// fetch current user info
let response;
try {
response = await userService.currentUser(cookie);
} catch (error) {
manualLogger.error("Failed to fetch current user:", error);
throw error;
console.log("caught?");
handleError(error, {
errorType: "unauthorized",
message: "Failed to authenticate user",
component: "authentication",
operation: "fetch-current-user",
extraContext: {
userId,
workspaceSlug,
},
throw: true,
});
}
if (response.id !== userId) {
throw Error("Authentication failed: Token doesn't match the current user.");
handleError(null, {
errorType: "unauthorized",
message: "Authentication failed: Token doesn't match the current user.",
component: "authentication",
operation: "validate-user",
extraContext: {
userId,
workspaceSlug,
},
throw: true,
});
}
return {
+24 -72
View File
@@ -1,112 +1,64 @@
// helpers
import {
getAllDocumentFormatsFromBinaryData,
getBinaryDataFromHTMLString,
} from "@/core/helpers/page.js";
import { getAllDocumentFormatsFromBinaryData, getBinaryDataFromHTMLString } from "@/core/helpers/page";
// services
import { PageService } from "@/core/services/page.service.js";
import { manualLogger } from "../helpers/logger.js";
import { PageService } from "@/services/page.service";
const pageService = new PageService();
export const updatePageDescription = async (
params: URLSearchParams,
pageId: string,
updatedDescription: Uint8Array,
cookie: string | undefined,
cookie: string | undefined
) => {
if (!(updatedDescription instanceof Uint8Array)) {
throw new Error(
"Invalid updatedDescription: must be an instance of Uint8Array",
);
throw new Error("Invalid updatedDescription: must be an instance of Uint8Array");
}
const workspaceSlug = params.get("workspaceSlug")?.toString();
const projectId = params.get("projectId")?.toString();
if (!workspaceSlug || !projectId || !cookie) return;
const { contentBinaryEncoded, contentHTML, contentJSON } =
getAllDocumentFormatsFromBinaryData(updatedDescription);
try {
const payload = {
description_binary: contentBinaryEncoded,
description_html: contentHTML,
description: contentJSON,
};
const { contentBinaryEncoded, contentHTML, contentJSON } = getAllDocumentFormatsFromBinaryData(updatedDescription);
const payload = {
description_binary: contentBinaryEncoded,
description_html: contentHTML,
description: contentJSON,
};
await pageService.updateDescription(
workspaceSlug,
projectId,
pageId,
payload,
cookie,
);
} catch (error) {
manualLogger.error("Update error:", error);
throw error;
}
await pageService.updateDescription(workspaceSlug, projectId, pageId, payload, cookie);
};
const fetchDescriptionHTMLAndTransform = async (
workspaceSlug: string,
projectId: string,
pageId: string,
cookie: string,
cookie: string
) => {
if (!workspaceSlug || !projectId || !cookie) return;
try {
const pageDetails = await pageService.fetchDetails(
workspaceSlug,
projectId,
pageId,
cookie,
);
const { contentBinary } = getBinaryDataFromHTMLString(
pageDetails.description_html ?? "<p></p>",
);
return contentBinary;
} catch (error) {
manualLogger.error(
"Error while transforming from HTML to Uint8Array",
error,
);
throw error;
}
const pageDetails = await pageService.fetchDetails(workspaceSlug, projectId, pageId, cookie);
const { contentBinary } = getBinaryDataFromHTMLString(pageDetails.description_html ?? "<p></p>");
return contentBinary;
};
export const fetchPageDescriptionBinary = async (
params: URLSearchParams,
pageId: string,
cookie: string | undefined,
cookie: string | undefined
) => {
const workspaceSlug = params.get("workspaceSlug")?.toString();
const projectId = params.get("projectId")?.toString();
if (!workspaceSlug || !projectId || !cookie) return null;
try {
const response = await pageService.fetchDescriptionBinary(
workspaceSlug,
projectId,
pageId,
cookie,
);
const binaryData = new Uint8Array(response);
const response = await pageService.fetchDescriptionBinary(workspaceSlug, projectId, pageId, cookie);
const binaryData = new Uint8Array(response);
if (binaryData.byteLength === 0) {
const binary = await fetchDescriptionHTMLAndTransform(
workspaceSlug,
projectId,
pageId,
cookie,
);
if (binary) {
return binary;
}
if (binaryData.byteLength === 0) {
const binary = await fetchDescriptionHTMLAndTransform(workspaceSlug, projectId, pageId, cookie);
if (binary) {
return binary;
}
return binaryData;
} catch (error) {
manualLogger.error("Fetch error:", error);
throw error;
}
return binaryData;
};
+62
View File
@@ -0,0 +1,62 @@
import { Redis } from "ioredis";
import { logger } from "@plane/logger";
let redisClient: Redis | null = null;
export async function initializeRedis(): Promise<Redis> {
const redisUrl = getRedisUrl();
if (!redisUrl) {
logger.error("Redis URL is not configured. Please set REDIS_URL environment variable.");
process.exit(1);
}
try {
redisClient = new Redis(redisUrl);
redisClient.on("error", (error) => {
logger.error("Redis connection error:", error);
process.exit(1);
});
// Wait for the connection to be ready
await new Promise<void>((resolve, reject) => {
redisClient!.on("ready", () => {
logger.info("Redis connection established successfully");
resolve();
});
redisClient!.on("error", (error) => {
reject(error);
});
});
return redisClient;
} catch (error) {
logger.error("Failed to initialize Redis:", error);
process.exit(1);
}
}
export function getRedisClient(): Redis {
if (!redisClient) {
throw new Error("Redis client not initialized. Call initializeRedis() first.");
}
return redisClient;
}
function getRedisUrl() {
const redisUrl = process.env.REDIS_URL?.trim();
const redisHost = process.env.REDIS_HOST?.trim();
const redisPort = process.env.REDIS_PORT?.trim();
if (redisUrl) {
return redisUrl;
}
if (redisHost && redisPort && !Number.isNaN(Number(redisPort))) {
return `redis://${redisHost}:${redisPort}`;
}
return "";
}
-15
View File
@@ -1,15 +0,0 @@
export function getRedisUrl() {
const redisUrl = process.env.REDIS_URL?.trim();
const redisHost = process.env.REDIS_HOST?.trim();
const redisPort = process.env.REDIS_PORT?.trim();
if (redisUrl) {
return redisUrl;
}
if (redisHost && redisPort && !Number.isNaN(Number(redisPort))) {
return `redis://${redisHost}:${redisPort}`;
}
return "";
}
-78
View File
@@ -1,78 +0,0 @@
// types
import { TPage } from "@plane/types";
// services
import { API_BASE_URL, APIService } from "@/core/services/api.service.js";
export class PageService extends APIService {
constructor() {
super(API_BASE_URL);
}
async fetchDetails(
workspaceSlug: string,
projectId: string,
pageId: string,
cookie: string
): Promise<TPage> {
return this.get(
`/api/workspaces/${workspaceSlug}/projects/${projectId}/pages/${pageId}/`,
{
headers: {
Cookie: cookie,
},
}
)
.then((response) => response?.data)
.catch((error) => {
throw error?.response?.data;
});
}
async fetchDescriptionBinary(
workspaceSlug: string,
projectId: string,
pageId: string,
cookie: string
): Promise<any> {
return this.get(
`/api/workspaces/${workspaceSlug}/projects/${projectId}/pages/${pageId}/description/`,
{
headers: {
"Content-Type": "application/octet-stream",
Cookie: cookie,
},
responseType: "arraybuffer",
}
)
.then((response) => response?.data)
.catch((error) => {
throw error?.response?.data;
});
}
async updateDescription(
workspaceSlug: string,
projectId: string,
pageId: string,
data: {
description_binary: string;
description_html: string;
description: object;
},
cookie: string
): Promise<any> {
return this.patch(
`/api/workspaces/${workspaceSlug}/projects/${projectId}/pages/${pageId}/description/`,
data,
{
headers: {
Cookie: cookie,
},
}
)
.then((response) => response?.data)
.catch((error) => {
throw error;
});
}
}
+6 -1
View File
@@ -1,10 +1,15 @@
// types
import { TAdditionalDocumentTypes } from "@/plane-live/types/common.js";
import { TAdditionalDocumentTypes } from "@/plane-live/types/common";
export type TDocumentTypes = "project_page" | TAdditionalDocumentTypes;
export type HocusPocusServerContext = {
cookie: string;
projectId: string;
workspaceSlug: string;
documentType: TDocumentTypes;
userId: string;
agentId: string;
};
export type TConvertDocumentRequestBody = {
+62
View File
@@ -0,0 +1,62 @@
import { HocusPocusServerContext, TDocumentTypes } from "@/core/types/common";
/**
* Parameters for document fetch operations
*/
export interface DocumentFetchParams {
context: HocusPocusServerContext;
pageId: string;
params: URLSearchParams;
}
/**
* Parameters for document store operations
*/
export interface DocumentStoreParams {
context: HocusPocusServerContext;
pageId: string;
state: any;
params: URLSearchParams;
}
/**
* Interface defining a document handler
*/
export interface DocumentHandler {
/**
* Fetch a document
*/
fetch: (params: DocumentFetchParams) => Promise<any>;
/**
* Store a document
*/
store: (params: DocumentStoreParams) => Promise<void>;
}
/**
* Handler context interface - extend this to add new criteria for handler selection
*/
export interface HandlerContext {
documentType?: TDocumentTypes;
agentId?: string;
}
/**
* Handler selector function type - determines if a handler should be used based on context
*/
export type HandlerSelector = (context: HandlerContext) => boolean;
/**
* Handler definition combining a selector and implementation
*/
export interface HandlerDefinition {
selector: HandlerSelector;
handler: DocumentHandler;
priority: number; // Higher number means higher priority
}
/**
* Type for a handler registration function
*/
export type RegisterHandler = (definition: HandlerDefinition) => void;
+1
View File
@@ -0,0 +1 @@
export * from "./register";
+5
View File
@@ -0,0 +1,5 @@
import { registerProjectPageHandler } from "@/ce/document-types/project-page";
export function initializeDocumentHandlers() {
registerProjectPageHandler();
}
+2 -1
View File
@@ -1 +1,2 @@
export * from "../../ce/lib/authentication.js"
export * from "../../core/lib/authentication";
+2 -1
View File
@@ -1 +1,2 @@
export * from "../../ce/lib/fetch-document.js"
export * from "../../ce/lib/fetch-document";
+2 -1
View File
@@ -1 +1,2 @@
export * from "../../ce/lib/update-document.js"
export * from "../../ce/lib/update-document";
+1 -1
View File
@@ -1 +1 @@
export * from "../../ce/types/common.js"
export * from "../../ce/types/common"
+42
View File
@@ -0,0 +1,42 @@
import * as dotenv from "dotenv";
import { z } from "zod";
// Load environment variables from .env file
dotenv.config();
// Define environment schema with validation
const envSchema = z.object({
// Server configuration
NODE_ENV: z.enum(["development", "test", "production"]).default("development"),
PORT: z.string().default("3000").transform(Number),
LIVE_BASE_PATH: z.string().default("/live"),
// CORS configuration
CORS_ALLOWED_ORIGINS: z.string().default("*"),
// Compression options
COMPRESSION_LEVEL: z.string().default("6").transform(Number),
COMPRESSION_THRESHOLD: z.string().default("5000").transform(Number),
// Hocuspocus server configuration
HOCUSPOCUS_URL: z.string().optional(),
HOCUSPOCUS_USERNAME: z.string().optional(),
HOCUSPOCUS_PASSWORD: z.string().optional(),
// Graceful termination timeout
SHUTDOWN_TIMEOUT: z.string().default("10000").transform(Number),
});
// Validate the environment variables
function validateEnv() {
const result = envSchema.safeParse(process.env);
if (!result.success) {
console.error("❌ Invalid environment variables:", JSON.stringify(result.error.format(), null, 4));
process.exit(1);
}
return result.data;
}
// Export the validated environment
export const env = validateEnv();
+145
View File
@@ -0,0 +1,145 @@
import { DocumentCollaborativeEvents, TDocumentEventsServer } from "@plane/editor/lib";
import { Logger } from "@hocuspocus/extension-logger";
import { Database } from "@hocuspocus/extension-database";
import { Redis } from "@hocuspocus/extension-redis";
import { Hocuspocus } from "@hocuspocus/server";
import { v4 as uuidv4 } from "uuid";
import { getRedisClient } from "./redis";
import { UserService } from "./services/user.service";
// import { handleError } from "@/core/helpers/error-handling/error-factory";
// import { TDocumentTypes } from "@/core/types/common";
// import { handleAuthentication } from "@/core/lib/authentication";
// import { IncomingHttpHeaders } from "http";
export const createHocusPocus = () => {
const serverName = process.env.HOSTNAME || uuidv4();
return new Hocuspocus({
name: serverName,
onAuthenticate: onAuthenticate(),
onStateless: onStateless(),
extensions: [
new Logger(),
new Database({
fetch: handleDataFetch,
store: handleDataStore,
}),
new Redis({ redis: getRedisClient() }),
],
debounce: 1000,
});
};
const validateToken = async (token: string | undefined) => {
try {
if (!token) {
throw new Error("Token not provided");
}
const userService = new UserService();
const response = await userService.currentUser(token);
return response;
} catch (error) {
return null;
}
};
const onAuthenticate = () => {
return async ({ token }: { token: string | undefined }) => {
const user = await validateToken(token);
if (!user) {
throw new Error("Invalid token");
}
return user;
};
};
const onStateless = () => {
return async ({ payload, document }: any) => {
// broadcast the client event (derived from the server event) to all the clients so that they can update their state
const response = DocumentCollaborativeEvents[payload as TDocumentEventsServer].client;
if (response) {
document.broadcastStateless(response);
}
};
};
const handleDataFetch = async (data: any) => {
try {
const { context, documentName, requestParameters } = data;
console.log("handleDataFetch", context);
console.log("handleDataFetch", documentName);
console.log("handleDataFetch", requestParameters);
return documentName; // TODO: remove this once the API integration is done
// fetch the data using page service
// const pageService = new PageService();
// const page = await pageService.getPage(documentName);
// return page;
} catch (error) {
console.error("handleDataFetch", error);
throw error;
}
};
const handleDataStore = async (data: any) => {
try {
const { context, documentName, requestParameters } = data;
console.log("handleDataStore", context);
console.log("handleDataStore", documentName);
console.log("handleDataStore", requestParameters);
return documentName; // TODO: remove this once the API integration is done
// store the data using page service
// const pageService = new PageService();
// const page = await pageService.updatePage(documentName, requestParameters);
// return page;
} catch (error) {
console.error("handleDataStore", error);
throw error;
}
};
// async ({
// token,
// requestParameters,
// requestHeaders,
// }: {
// token: string;
// requestParameters: URLSearchParams;
// requestHeaders: IncomingHttpHeaders;
// }) => {
// let cookie: string | undefined = undefined;
// let userId: string | undefined = undefined;
// try {
// const parsedToken = JSON.parse(token) as { id: string; cookie: string };
// userId = parsedToken.id;
// cookie = parsedToken.cookie;
// } catch (error) {
// console.error("Token parsing failed, using request headers:", error);
// } finally {
// if (!cookie) {
// cookie = requestHeaders.cookie?.toString();
// }
// }
// if (!cookie || !userId) {
// handleError(null, {
// errorType: "unauthorized",
// message: "Credentials not provided",
// component: "hocuspocus",
// operation: "authenticate",
// extraContext: { tokenProvided: !!token },
// throw: true,
// });
// }
// const documentType = requestParameters.get("documentType")?.toString() as TDocumentTypes;
// const workspaceSlug = requestParameters.get("workspaceSlug")?.toString() as string;
// return await handleAuthentication({
// cookie,
// userId,
// workspaceSlug,
// });
// }
@@ -0,0 +1,278 @@
import { AppError, HttpStatusCode, ErrorCategory } from "./error-handler";
/**
* Map of error types to their corresponding factory functions
* This ensures that error types and their implementations stay in sync
*/
interface ErrorFactory {
statusCode: number;
category: ErrorCategory;
defaultMessage: string;
createError: (message?: string, context?: Record<string, any>) => AppError;
}
const ERROR_FACTORIES = {
"bad-request": {
statusCode: HttpStatusCode.BAD_REQUEST,
category: ErrorCategory.OPERATIONAL,
defaultMessage: "Bad Request",
createError: (message = "Bad Request", context?) =>
new AppError(message, HttpStatusCode.BAD_REQUEST, ErrorCategory.OPERATIONAL, context),
},
unauthorized: {
statusCode: HttpStatusCode.UNAUTHORIZED,
category: ErrorCategory.OPERATIONAL,
defaultMessage: "Unauthorized",
createError: (message = "Unauthorized", context?) =>
new AppError(message, HttpStatusCode.UNAUTHORIZED, ErrorCategory.OPERATIONAL, context),
},
forbidden: {
statusCode: HttpStatusCode.FORBIDDEN,
category: ErrorCategory.OPERATIONAL,
defaultMessage: "Forbidden",
createError: (message = "Forbidden", context?) =>
new AppError(message, HttpStatusCode.FORBIDDEN, ErrorCategory.OPERATIONAL, context),
},
"not-found": {
statusCode: HttpStatusCode.NOT_FOUND,
category: ErrorCategory.OPERATIONAL,
defaultMessage: "Resource not found",
createError: (message = "Resource not found", context?) =>
new AppError(message, HttpStatusCode.NOT_FOUND, ErrorCategory.OPERATIONAL, context),
},
conflict: {
statusCode: HttpStatusCode.CONFLICT,
category: ErrorCategory.OPERATIONAL,
defaultMessage: "Resource conflict",
createError: (message = "Resource conflict", context?) =>
new AppError(message, HttpStatusCode.CONFLICT, ErrorCategory.OPERATIONAL, context),
},
"unprocessable-entity": {
statusCode: HttpStatusCode.UNPROCESSABLE_ENTITY,
category: ErrorCategory.OPERATIONAL,
defaultMessage: "Unprocessable Entity",
createError: (message = "Unprocessable Entity", context?) =>
new AppError(message, HttpStatusCode.UNPROCESSABLE_ENTITY, ErrorCategory.OPERATIONAL, context),
},
"too-many-requests": {
statusCode: HttpStatusCode.TOO_MANY_REQUESTS,
category: ErrorCategory.OPERATIONAL,
defaultMessage: "Too many requests",
createError: (message = "Too many requests", context?) =>
new AppError(message, HttpStatusCode.TOO_MANY_REQUESTS, ErrorCategory.OPERATIONAL, context),
},
internal: {
statusCode: HttpStatusCode.INTERNAL_SERVER,
category: ErrorCategory.PROGRAMMING,
defaultMessage: "Internal Server Error",
createError: (message = "Internal Server Error", context?) =>
new AppError(message, HttpStatusCode.INTERNAL_SERVER, ErrorCategory.PROGRAMMING, context),
},
"service-unavailable": {
statusCode: HttpStatusCode.SERVICE_UNAVAILABLE,
category: ErrorCategory.SYSTEM,
defaultMessage: "Service Unavailable",
createError: (message = "Service Unavailable", context?) =>
new AppError(message, HttpStatusCode.SERVICE_UNAVAILABLE, ErrorCategory.SYSTEM, context),
},
fatal: {
statusCode: HttpStatusCode.INTERNAL_SERVER,
category: ErrorCategory.FATAL,
defaultMessage: "Fatal Error",
createError: (message = "Fatal Error", context?) =>
new AppError(message, HttpStatusCode.INTERNAL_SERVER, ErrorCategory.FATAL, context),
},
} satisfies Record<string, ErrorFactory>;
// Create the type from the keys of the error factories map
export type ErrorType = keyof typeof ERROR_FACTORIES;
// -------------------------------------------------------------------------
// Primary public API - Recommended for most use cases
// -------------------------------------------------------------------------
/**
* Base options for handleError function
*/
type BaseErrorHandlerOptions = {
// Error classification options
errorType?: ErrorType;
message?: string;
// Context information
component: string;
operation: string;
extraContext?: Record<string, any>;
// Behavior options
rethrowIfAppError?: boolean;
};
/**
* Options for throwing variant of handleError - discriminated by throw: true
*/
export type ThrowingOptions = BaseErrorHandlerOptions & {
throw: true;
};
/**
* Options for non-throwing variant of handleError - default behavior
*/
export type NonThrowingOptions = BaseErrorHandlerOptions;
/**
* Unified error handler that encapsulates common error handling patterns
*
* @param error The error to handle
* @param options Configuration options with throw: true to throw the error instead of returning it
* @returns Never returns - always throws
* @example
* // Throwing version
* handleError(error, {
* errorType: 'not-found',
* component: 'user-service',
* operation: 'getUserById',
* throw: true
* });
*/
export function handleError(error: unknown, options: ThrowingOptions): never;
/**
* Unified error handler that encapsulates common error handling patterns
*
* @param error The error to handle
* @param options Configuration options (non-throwing by default)
* @returns The AppError instance
* @example
* // Non-throwing version (default)
* const appError = handleError(error, {
* errorType: 'not-found',
* component: 'user-service',
* operation: 'getUserById'
* });
* return { error: appError.output() };
*/
// eslint-disable-next-line no-redeclare
export function handleError(error: unknown, options: NonThrowingOptions): AppError;
/**
* Implementation of handleError that handles both throwing and non-throwing cases
*/
// eslint-disable-next-line no-redeclare
export function handleError(error: unknown, options: ThrowingOptions | NonThrowingOptions): AppError | never {
// Only throw if throw is explicitly true
const shouldThrow = (options as ThrowingOptions).throw === true;
// If the error is already an AppError and we want to rethrow it as is
if (options.rethrowIfAppError !== false && error instanceof AppError) {
if (shouldThrow) {
throw error;
}
return error;
}
// Format the error message
const errorMessage = options.message
? error instanceof Error
? `${options.message}: ${error.message}`
: error
? `${options.message}: ${String(error)}`
: options.message
: error instanceof Error
? error.message
: error
? String(error)
: "Unknown error occurred";
// Build context object
const context = {
component: options.component,
operation: options.operation,
originalError: error,
...(options.extraContext || {}),
};
// Create the appropriate error type using our factory map
const errorType = options.errorType || "internal";
const factory = ERROR_FACTORIES[errorType];
if (!factory) {
// If no factory found, default to internal error
return ERROR_FACTORIES.internal.createError(errorMessage, context);
}
// Create the error with the factory
const appError = factory.createError(errorMessage, context);
// If we should throw, do so now
if (shouldThrow) {
throw appError;
}
return appError;
}
/**
* Utility function to convert errors or enhance existing AppErrors
*/
export const convertError = (
error: Error,
options?: {
statusCode?: number;
message?: string;
category?: ErrorCategory;
context?: Record<string, any>;
}
): AppError => {
if (error instanceof AppError) {
// If it's already an AppError and no overrides, return as is
if (!options?.statusCode && !options?.message && !options?.category) {
return error;
}
// Create a new AppError with the original as context
return new AppError(
options?.message || error.message,
options?.statusCode || error.status,
options?.category || error.category,
{
...(error.context || {}),
...(options?.context || {}),
originalError: error,
}
);
}
// Determine the appropriate error type based on status code
let errorType: ErrorType = "internal";
if (options?.statusCode) {
// Find the error type that matches the status code
const entry = Object.entries(ERROR_FACTORIES).find(([_, factory]) => factory.statusCode === options.statusCode);
if (entry) {
errorType = entry[0] as ErrorType;
}
}
// Return a new AppError using the factory
return handleError(error, {
errorType: errorType,
message: options?.message,
component: options?.context?.component || "unknown",
operation: options?.context?.operation || "convert-error",
extraContext: options?.context,
});
};
/**
* Check if an error is an AppError
*/
export const isAppError = (err: any, statusCode?: number): boolean => {
return err instanceof AppError && (!statusCode || err.status === statusCode);
};
// Export only the public API
export default {
handleError,
convertError,
isAppError,
};
@@ -0,0 +1,384 @@
import { ErrorRequestHandler, Request, Response, NextFunction } from "express";
import { env } from "@/env";
import { logger } from "@plane/logger";
import { handleError } from "./error-factory";
import { ErrorContext, reportError } from "./error-reporting";
import { manualLogger } from "../../core/helpers/logger";
/**
* HTTP Status Codes
*/
export enum HttpStatusCode {
// 2xx Success
OK = 200,
CREATED = 201,
ACCEPTED = 202,
NO_CONTENT = 204,
// 4xx Client Errors
BAD_REQUEST = 400,
UNAUTHORIZED = 401,
FORBIDDEN = 403,
NOT_FOUND = 404,
METHOD_NOT_ALLOWED = 405,
CONFLICT = 409,
GONE = 410,
UNPROCESSABLE_ENTITY = 422,
TOO_MANY_REQUESTS = 429,
// 5xx Server Errors
INTERNAL_SERVER = 500,
NOT_IMPLEMENTED = 501,
BAD_GATEWAY = 502,
SERVICE_UNAVAILABLE = 503,
GATEWAY_TIMEOUT = 504,
}
/**
* Error categories to classify errors
*/
export enum ErrorCategory {
OPERATIONAL = "operational", // Expected errors that are part of normal operation (e.g. validation failures)
PROGRAMMING = "programming", // Unexpected errors that indicate bugs (e.g. null references)
SYSTEM = "system", // System errors (e.g. out of memory, connection failures)
FATAL = "fatal", // Severe errors that should crash the app (e.g. unrecoverable state)
}
/**
* Base Application Error Class
* All custom errors extend this class
*/
export class AppError extends Error {
readonly status: number;
readonly category: ErrorCategory;
readonly context?: Record<string, any>;
readonly isOperational: boolean; // Kept for backward compatibility
constructor(
message: string,
status: number = HttpStatusCode.INTERNAL_SERVER,
category: ErrorCategory = ErrorCategory.PROGRAMMING,
context?: Record<string, any>
) {
super(message);
// Set error properties
this.name = this.constructor.name;
this.status = status;
this.category = category;
this.isOperational = category === ErrorCategory.OPERATIONAL;
this.context = context;
// Capture stack trace, excluding the constructor call from the stack
Error.captureStackTrace(this, this.constructor);
// Automatically report the error (unless it's being constructed by the error utilities)
if (!context?.skipReporting) {
this.report();
}
}
/**
* Creates a formatted representation of the error
*/
output() {
return {
statusCode: this.status,
payload: {
statusCode: this.status,
error: this.getErrorName(),
message: this.message,
category: this.category,
},
headers: {},
};
}
/**
* Gets a descriptive name for the error based on status code
*/
private getErrorName(): string {
const statusCodes: Record<number, string> = {
400: "Bad Request",
401: "Unauthorized",
403: "Forbidden",
404: "Not Found",
405: "Method Not Allowed",
409: "Conflict",
410: "Gone",
422: "Unprocessable Entity",
429: "Too Many Requests",
500: "Internal Server Error",
501: "Not Implemented",
502: "Bad Gateway",
503: "Service Unavailable",
504: "Gateway Timeout",
};
return statusCodes[this.status] || "Unknown Error";
}
/**
* Reports the error to logging and monitoring systems
*/
private report(): void {
// Different logging based on error category
if (this.category === ErrorCategory.OPERATIONAL) {
manualLogger.error(`Operational error: ${this.message}`, {
errorName: this.name,
errorStatus: this.status,
errorCategory: this.category,
context: this.context,
});
} else if (this.category === ErrorCategory.FATAL) {
manualLogger.error(`FATAL error: ${this.message}`, {
errorName: this.name,
errorStatus: this.status,
errorCategory: this.category,
stack: this.stack,
context: this.context,
});
} else {
manualLogger.error(`${this.category} error: ${this.message}`, {
errorName: this.name,
errorStatus: this.status,
errorCategory: this.category,
stack: this.stack,
context: this.context,
});
}
}
}
export class FatalError extends AppError {
constructor(message: string, context?: Record<string, any>) {
super(message, HttpStatusCode.INTERNAL_SERVER, ErrorCategory.FATAL, context);
}
}
/**
* Main Express error handler middleware
*/
export const errorHandler: ErrorRequestHandler = (err, req, res, next) => {
// Already sent response, let default Express error handler deal with it
if (res.headersSent) {
return next(err);
}
// Convert to AppError if it's not already one
const error = handleError(err, {
component: "express",
operation: "error-handler",
extraContext: {
originalError: err,
url: req.originalUrl,
method: req.method,
},
});
// Normalize status code
const statusCode = error.status;
// Set the response status
res.status(statusCode);
// Set any custom headers if provided in the error object
if (err.headers && typeof err.headers === "object") {
Object.entries(err.headers).forEach(([key, value]) => {
res.set(key, value as string);
});
}
// Prepare error response
const errorResponse: {
error: {
message: string;
status: number;
stack?: string;
};
} = {
error: {
message:
error.category === ErrorCategory.OPERATIONAL || env.NODE_ENV !== "production"
? error.message
: "An unexpected error occurred",
status: statusCode,
},
};
// Add stack trace in non-production environments
if (env.NODE_ENV !== "production") {
errorResponse.error.stack = error.stack;
}
// Send the response
res.json(errorResponse);
// For fatal errors, log but NEVER terminate the app
if (error.category === ErrorCategory.FATAL) {
logger.error(`FATAL ERROR OCCURRED BUT APP WILL CONTINUE RUNNING: ${error.message}`);
}
};
export const asyncHandler = (fn: Function) => {
return (req: any, res: any, next: any) => {
Promise.resolve(fn(req, res, next)).catch((error) => {
// Convert to AppError if needed and pass to Express error middleware
const appError = handleError(error, {
errorType: "internal",
component: "express",
operation: "route-handler",
extraContext: {
url: req.originalUrl,
method: req.method,
body: req.body,
query: req.query,
params: req.params,
},
});
next(appError);
});
};
};
export interface CatchAsyncOptions<T, E = Error> {
/** Default value to return in case of error, null by default */
defaultValue?: T | null;
/** Whether to report non-AppErrors automatically */
reportErrors?: boolean;
/** Whether to rethrow the error after handling it */
rethrow?: boolean;
/** Custom error transformer function */
transformError?: (error: unknown) => E;
/** Custom error handler function that runs before standard handling */
onError?: (error: unknown) => void | Promise<void>;
/** Custom handler for specific error types */
errorHandlers?: {
[key: string]: (error: any) => T | null | Promise<T>;
};
}
export const catchAsync = <T, E = Error>(
fn: () => Promise<T>,
context?: ErrorContext,
options: CatchAsyncOptions<T, E> = {}
): (() => Promise<T | null>) => {
const { defaultValue = null, onError, rethrow = false } = options;
return async () => {
try {
return await fn();
} catch (error) {
// Apply custom error handler if provided
if (onError) {
await Promise.resolve(onError(error));
}
reportError(error, context);
if (error instanceof AppError) {
error.context;
}
if (rethrow) {
// Use handleError to ensure consistent error handling when rethrowing
handleError(error, {
component: context?.extra?.component || "unknown",
operation: context?.extra?.operation || "unknown",
extraContext: {
...context,
...(error instanceof AppError ? error.context : {}),
originalError: error,
},
throw: true,
});
}
return defaultValue;
}
};
};
/**
* Set up global error handlers for uncaught exceptions and unhandled rejections
* @param gracefulTerminationHandler Function to call for graceful termination
*/
export const setupGlobalErrorHandlers = (gracefulTerminationHandler: () => Promise<void>): void => {
// Handle promise rejections
process.on("unhandledRejection", (reason: unknown) => {
logger.error("Unhandled Promise Rejection", { reason });
// Convert to AppError and handle
const appError = handleError(reason, {
errorType: "internal",
message: reason instanceof Error ? reason.message : String(reason),
component: "process",
operation: "unhandledRejection",
extraContext: { source: "unhandledRejection" },
});
// Log the error but never terminate
logger.error(`Unhandled rejection caught and contained: ${appError.message}`);
});
// Handle exceptions
process.on("uncaughtException", (error: Error) => {
logger.error("Uncaught Exception", {
error: error.message,
stack: error.stack,
});
// Convert to AppError if needed
const appError = handleError(error, {
errorType: "internal",
component: "process",
operation: "uncaughtException",
extraContext: {
source: "uncaughtException",
},
});
// Log the error but never terminate
logger.warn(`Uncaught exception contained: ${appError.message}`);
});
// Handle termination signals
process.on("SIGTERM", () => {
logger.info("SIGTERM received. Starting graceful termination...");
gracefulTerminationHandler();
});
process.on("SIGINT", () => {
logger.info("SIGINT received. Starting graceful termination...");
gracefulTerminationHandler();
});
};
/**
* Configure error handling middleware for the Express app
* @param app Express application instance
*/
export function configureErrorHandlers(app: any): void {
// Global error handling middleware
app.use(errorHandler);
// 404 handler must be last
app.use((_req: Request, _res: Response, next: NextFunction) => {
next(
handleError(null, {
errorType: "not-found",
message: "Resource not found",
component: "express",
operation: "route-handler",
extraContext: { path: _req.path },
throw: true,
})
);
});
}
@@ -0,0 +1,45 @@
import { AppError, ErrorCategory } from "./error-handler";
import { logger } from "@plane/logger";
import { handleError } from "./error-factory";
export interface ErrorContext {
url?: string;
method?: string;
body?: any;
query?: any;
params?: any;
extra?: Record<string, any>;
}
/**
* Utility function to report errors that aren't instances of AppError
* AppError instances automatically report themselves on creation
* Only use this for external errors that don't use our error system
*/
export const reportError = (error: Error | unknown, context?: ErrorContext): void => {
if (error instanceof AppError) {
// if it's an app error, don't report it as it's already been reported
return;
}
logger.error(`External error: ${error instanceof Error ? error.stack || error.message : String(error)}`, {
error,
context,
});
};
export const handleFatalError = (error: Error | unknown, context?: ErrorContext): void => {
// Convert to fatal AppError
const fatalError = handleError(error, {
errorType: "fatal",
message: error instanceof Error ? error.message : String(error),
component: context?.extra?.component || "system",
operation: context?.extra?.operation || "fatal-error-handler",
extraContext: {
...context,
originalError: error,
},
});
process.emit("uncaughtException", fatalError);
};
+62
View File
@@ -0,0 +1,62 @@
import { Redis } from "ioredis";
import { logger } from "@plane/logger";
let redisClient: Redis | null = null;
export async function initializeRedis(): Promise<Redis> {
const redisUrl = getRedisUrl();
if (!redisUrl) {
logger.error("Redis URL is not configured. Please set REDIS_URL environment variable.");
process.exit(1);
}
try {
redisClient = new Redis(redisUrl);
redisClient.on("error", (error) => {
logger.error("Redis connection error:", error);
process.exit(1);
});
// Wait for the connection to be ready
await new Promise<void>((resolve, reject) => {
redisClient!.on("ready", () => {
logger.info("Redis connection established successfully");
resolve();
});
redisClient!.on("error", (error) => {
reject(error);
});
});
return redisClient;
} catch (error) {
logger.error("Failed to initialize Redis:", error);
process.exit(1);
}
}
export function getRedisClient(): Redis {
if (!redisClient) {
throw new Error("Redis client not initialized. Call initializeRedis() first.");
}
return redisClient;
}
function getRedisUrl() {
const redisUrl = process.env.REDIS_URL?.trim();
const redisHost = process.env.REDIS_HOST?.trim();
const redisPort = process.env.REDIS_PORT?.trim();
if (redisUrl) {
return redisUrl;
}
if (redisHost && redisPort && !Number.isNaN(Number(redisPort))) {
return `redis://${redisHost}:${redisPort}`;
}
return "";
}
+106 -122
View File
@@ -1,135 +1,119 @@
import compression from "compression";
import cookieParser from "cookie-parser";
import cors from "cors";
import express, { Application } from "express";
import expressWs from "express-ws";
import express from "express";
import helmet from "helmet";
// hocuspocus server
import { getHocusPocusServer } from "@/core/hocuspocus-server.js";
// helpers
import { convertHTMLDocumentToAllFormats } from "@/core/helpers/convert-document.js";
import { logger, manualLogger } from "@/core/helpers/logger.js";
import { errorHandler } from "@/core/helpers/error-handler.js";
// types
import { TConvertDocumentRequestBody } from "@/core/types/common.js";
import path from "path";
import { Hocuspocus } from "@hocuspocus/server";
// controllers
const app: any = express();
expressWs(app);
import { initializeRedis } from "./core/lib/redis-manager";
import { logger } from "@plane/logger";
import { createHocusPocus } from "./hocuspocus";
app.set("port", process.env.PORT || 3000);
import { registerWebSocketController, registerController } from "@plane/decorators";
import { REST_CONTROLLERS, WEBSOCKET_CONTROLLERS } from "./controllers";
// Security middleware
app.use(helmet());
export default class Server {
app: Application;
PORT: number;
BASE_PATH: string;
CORS_ALLOWED_ORIGINS: string;
hocuspocusServer: Hocuspocus | null = null;
private httpServer: any;
// Middleware for response compression
app.use(
compression({
level: 6,
threshold: 5 * 1000,
})
);
// Logging middleware
app.use(logger);
// Body parsing middleware
app.use(express.json());
app.use(express.urlencoded({ extended: true }));
// cors middleware
app.use(cors());
const router = express.Router();
const HocusPocusServer = await getHocusPocusServer().catch((err) => {
manualLogger.error("Failed to initialize HocusPocusServer:", err);
process.exit(1);
});
router.get("/health", (_req, res) => {
res.status(200).json({ status: "OK" });
});
router.ws("/collaboration", (ws, req) => {
try {
HocusPocusServer.handleConnection(ws, req);
} catch (err) {
manualLogger.error("WebSocket connection error:", err);
ws.close();
constructor() {
this.PORT = parseInt(process.env.PORT || "3000");
this.BASE_PATH = process.env.LIVE_BASE_PATH || "/";
this.CORS_ALLOWED_ORIGINS = process.env.CORS_ALLOWED_ORIGINS || "*";
// Initialize express app
this.app = express();
expressWs(this.app as any);
// Security middleware
this.app.use(helmet());
// cors
this.setupCors();
// Cookie parsing
this.app.use(cookieParser());
// Body parsing middleware
this.app.use(express.json());
this.app.use(express.urlencoded({ extended: true }));
// static files
this.app.use(express.static(path.join(__dirname, "public")));
// setup redis
initializeRedis();
// setup hocuspocus server
this.hocuspocusServer = createHocusPocus();
// setup controllers
this.setupControllers();
}
});
router.post("/convert-document", (req, res) => {
const { description_html, variant } = req.body as TConvertDocumentRequestBody;
try {
if (description_html === undefined || variant === undefined) {
res.status(400).send({
message: "Missing required fields",
});
private setupCors() {
const origins = this.CORS_ALLOWED_ORIGINS.split(",").map((origin) => origin.trim());
this.app.use(
cors({
origin: origins,
credentials: true,
methods: ["GET", "POST", "PUT", "DELETE", "OPTIONS"],
allowedHeaders: ["Content-Type", "Authorization", "x-api-key"],
})
);
}
private setupControllers() {
const router = express.Router();
REST_CONTROLLERS.forEach((controller: any) => {
registerController(router, controller, [this.hocuspocusServer]);
});
WEBSOCKET_CONTROLLERS.forEach((controller: any) => {
registerWebSocketController(router, controller, [this.hocuspocusServer]);
});
this.app.use(this.BASE_PATH, router);
}
start() {
this.httpServer = this.app.listen(this.PORT, () => {
console.log(`Plane Live server has started at port ${this.PORT}`);
});
// Setup graceful shutdown
process.on("SIGTERM", () => this.shutdown("Received SIGTERM"));
process.on("SIGINT", () => this.shutdown("Received SIGINT"));
process.on("uncaughtException", (error) => {
logger.error("Uncaught exception:", error);
});
// Handle unhandled promise rejections - create AppError but DON'T terminate
process.on("unhandledRejection", (error) => {
logger.error("Unhandled rejection:", error);
});
}
private async shutdown(error: string): Promise<void> {
logger.info(`Initiating graceful shutdown: ${error}`);
if (!this.httpServer) {
logger.info("No HTTP server to close");
return;
}
const { description, description_binary } = convertHTMLDocumentToAllFormats({
document_html: description_html,
variant,
});
res.status(200).json({
description,
description_binary,
});
} catch (error) {
manualLogger.error("Error in /convert-document endpoint:", error);
res.status(500).send({
message: `Internal server error. ${error}`,
// Close all existing connections
this.httpServer.closeAllConnections?.();
// Close the server
return new Promise<void>((resolve) => {
this.httpServer.close((error: Error | undefined) => {
if (error) {
logger.error("Error closing HTTP server:", error);
} else {
logger.info("HTTP server closed successfully");
}
resolve();
});
});
}
});
app.use(process.env.LIVE_BASE_PATH || "/live", router);
app.use((_req, res) => {
res.status(404).send("Not Found");
});
app.use(errorHandler);
const liveServer = app.listen(app.get("port"), () => {
manualLogger.info(`Plane Live server has started at port ${app.get("port")}`);
});
const gracefulShutdown = async () => {
manualLogger.info("Starting graceful shutdown...");
try {
// Close the HocusPocus server WebSocket connections
await HocusPocusServer.destroy();
manualLogger.info("HocusPocus server WebSocket connections closed gracefully.");
// Close the Express server
liveServer.close(() => {
manualLogger.info("Express server closed gracefully.");
process.exit(1);
});
} catch (err) {
manualLogger.error("Error during shutdown:", err);
process.exit(1);
}
// Forcefully shut down after 10 seconds if not closed
setTimeout(() => {
manualLogger.error("Forcing shutdown...");
process.exit(1);
}, 10000);
};
// Graceful shutdown on unhandled rejection
process.on("unhandledRejection", (err: any) => {
manualLogger.info("Unhandled Rejection: ", err);
manualLogger.info(`UNHANDLED REJECTION! 💥 Shutting down...`);
gracefulShutdown();
});
// Graceful shutdown on uncaught exception
process.on("uncaughtException", (err: any) => {
manualLogger.info("Uncaught Exception: ", err);
manualLogger.info(`UNCAUGHT EXCEPTION! 💥 Shutting down...`);
gracefulShutdown();
});
}
+58
View File
@@ -0,0 +1,58 @@
// types
import { TPage } from "@plane/types";
// services
import { API_BASE_URL, APIService } from "@/services/api.service";
export class PageService extends APIService {
constructor() {
super(API_BASE_URL);
}
async fetchDetails(workspaceSlug: string, projectId: string, pageId: string, cookie: string): Promise<TPage> {
return this.get(`/api/workspaces/${workspaceSlug}/projects/${projectId}/pages/${pageId}/`, {
headers: {
Cookie: cookie,
},
})
.then((response) => response?.data)
.catch((error) => {
throw error?.response?.data;
});
}
async fetchDescriptionBinary(workspaceSlug: string, projectId: string, pageId: string, cookie: string): Promise<any> {
return this.get(`/api/workspaces/${workspaceSlug}/projects/${projectId}/pages/${pageId}/description/`, {
headers: {
"Content-Type": "application/octet-stream",
Cookie: cookie,
},
responseType: "arraybuffer",
})
.then((response) => response?.data)
.catch((error) => {
throw error?.response?.data;
});
}
async updateDescription(
workspaceSlug: string,
projectId: string,
pageId: string,
data: {
description_binary: string;
description_html: string;
description: object;
},
cookie: string
): Promise<any> {
return this.patch(`/api/workspaces/${workspaceSlug}/projects/${projectId}/pages/${pageId}/description/`, data, {
headers: {
Cookie: cookie,
},
})
.then((response) => response?.data)
.catch((error) => {
throw error;
});
}
}
@@ -1,7 +1,7 @@
// types
import type { IUser } from "@plane/types";
// services
import { API_BASE_URL, APIService } from "@/core/services/api.service.js";
import { API_BASE_URL, APIService } from "@/services/api.service";
export class UserService extends APIService {
constructor() {
+10
View File
@@ -0,0 +1,10 @@
import Server from "./server";
import { env } from "./env";
import { logger } from "@plane/logger";
// Log server startup details
logger.info(`Starting Plane Live server in ${env.NODE_ENV} environment`);
// Initialize and start the server
const server = new Server();
server.start();
+4 -2
View File
@@ -1,8 +1,8 @@
{
"extends": "@plane/typescript-config/base.json",
"compilerOptions": {
"module": "NodeNext",
"moduleResolution": "NodeNext",
"module": "ES2015",
"moduleResolution": "Bundler",
"lib": ["ES2015"],
"outDir": "./dist",
"rootDir": ".",
@@ -16,6 +16,8 @@
"skipLibCheck": true,
"sourceMap": true,
"inlineSources": true,
"experimentalDecorators": true,
"emitDecoratorMetadata": true,
"sourceRoot": "/"
},
"include": ["src/**/*.ts", "tsup.config.ts"],
+13 -9
View File
@@ -1,11 +1,15 @@
import { defineConfig, Options } from "tsup";
import { defineConfig } from "tsup";
export default defineConfig((options: Options) => ({
entry: ["src/server.ts"],
format: ["cjs", "esm"],
export default defineConfig({
entry: ["src/start.ts"],
format: ["esm", "cjs"],
dts: true,
clean: false,
external: ["react"],
injectStyle: true,
...options,
}));
splitting: false,
sourcemap: true,
minify: false,
target: "node18",
outDir: "dist",
env: {
NODE_ENV: process.env.NODE_ENV || "development",
},
});
+5 -4
View File
@@ -20,11 +20,12 @@ interface ControllerConstructor {
prototype: ControllerInstance;
}
export function registerControllers(
export function registerController(
router: Router,
Controller: ControllerConstructor,
dependencies: any[] = []
): void {
const instance = new Controller();
const instance = new Controller(...dependencies);
const baseRoute = Reflect.getMetadata("baseRoute", Controller) as string;
Object.getOwnPropertyNames(Controller.prototype).forEach((methodName) => {
@@ -33,14 +34,14 @@ export function registerControllers(
const method = Reflect.getMetadata(
"method",
instance,
methodName,
methodName
) as HttpMethod;
const route = Reflect.getMetadata("route", instance, methodName) as string;
const middlewares =
(Reflect.getMetadata(
"middlewares",
instance,
methodName,
methodName
) as RequestHandler[]) || [];
if (method && route) {
+2 -3
View File
@@ -2,8 +2,8 @@
export { Controller, Middleware } from "./rest";
export { Get, Post, Put, Patch, Delete } from "./rest";
export { WebSocket } from "./websocket";
export { registerControllers } from "./controller";
export { registerWebSocketControllers } from "./websocket-controller";
export { registerController } from "./controller";
export { registerWebSocketController } from "./websocket-controller";
// Also provide namespaced exports for better organization
import * as RestDecorators from "./rest";
@@ -12,4 +12,3 @@ import * as WebSocketDecorators from "./websocket";
// Named namespace exports
export const Rest = RestDecorators;
export const WebSocketNS = WebSocketDecorators;
+21 -56
View File
@@ -1,5 +1,3 @@
import { Router, Request } from "express";
import type { WebSocket } from "ws";
import "reflect-metadata";
interface ControllerInstance {
@@ -11,23 +9,19 @@ interface ControllerConstructor {
prototype: ControllerInstance;
}
export function registerWebSocketControllers(
router: Router,
export function registerWebSocketController(
router: any,
Controller: ControllerConstructor,
existingInstance?: ControllerInstance,
dependencies: any[] = []
): void {
const instance = existingInstance || new Controller();
const baseRoute = Reflect.getMetadata("baseRoute", Controller) as string;
const instance = new Controller(...dependencies);
const baseRoute = Reflect.getMetadata("baseRoute", Controller);
Object.getOwnPropertyNames(Controller.prototype).forEach((methodName) => {
if (methodName === "constructor") return; // Skip the constructor
const method = Reflect.getMetadata(
"method",
instance,
methodName,
) as string;
const route = Reflect.getMetadata("route", instance, methodName) as string;
const method = Reflect.getMetadata("method", instance, methodName);
const route = Reflect.getMetadata("route", instance, methodName);
if (method === "ws" && route) {
const handler = instance[methodName] as unknown;
@@ -36,50 +30,21 @@ export function registerWebSocketControllers(
typeof handler === "function" &&
typeof (router as any).ws === "function"
) {
(router as any).ws(
`${baseRoute}${route}`,
(ws: WebSocket, req: Request) => {
try {
handler.call(instance, ws, req);
} catch (error) {
console.error(
`WebSocket error in ${Controller.name}.${methodName}`,
error,
);
ws.close(
1011,
error instanceof Error
? error.message
: "Internal server error",
);
}
},
);
router.ws(`${baseRoute}${route}`, (ws: any, req: any) => {
try {
handler.call(instance, ws, req);
} catch (error) {
console.error(
`WebSocket error in ${Controller.name}.${methodName}`,
error
);
ws.close(
1011,
error instanceof Error ? error.message : "Internal server error"
);
}
});
}
}
});
}
/**
* Base controller class for WebSocket endpoints
*/
export abstract class BaseWebSocketController {
protected router: Router;
constructor() {
this.router = Router();
}
/**
* Get the base route for this controller
*/
protected getBaseRoute(): string {
return Reflect.getMetadata("baseRoute", this.constructor) || "";
}
/**
* Abstract method to handle WebSocket connections
* Implement this in your derived class
*/
abstract handleConnection(ws: WebSocket, req: Request): void;
}
+1 -1
View File
@@ -9,7 +9,7 @@ export function WebSocket(route: string): MethodDecorator {
return function (
target: object,
propertyKey: string | symbol,
descriptor: PropertyDescriptor,
descriptor: PropertyDescriptor
) {
Reflect.defineMetadata("method", "ws", target, propertyKey);
Reflect.defineMetadata("route", route, target, propertyKey);
+25 -3
View File
@@ -4,11 +4,18 @@
"license": "AGPL-3.0",
"description": "Logger shared across multiple apps internally",
"private": true,
"main": "./src/index.ts",
"types": "./src/index.ts",
"main": "./dist/index.js",
"module": "./dist/index.mjs",
"types": "./dist/index.d.ts",
"files": [
"dist/**"
],
"scripts": {
"build": "tsup",
"dev": "tsup --watch",
"lint": "eslint src --ext .ts,.tsx",
"lint:errors": "eslint src --ext .ts,.tsx --quiet"
"lint:errors": "eslint src --ext .ts,.tsx --quiet",
"clean": "rm -rf .turbo && rm -rf node_modules && rm -rf dist"
},
"dependencies": {
"winston": "^3.17.0",
@@ -17,6 +24,21 @@
"devDependencies": {
"@plane/eslint-config": "*",
"@types/node": "^22.5.4",
"tsup": "8.3.0",
"typescript": "^5.3.3"
},
"tsup": {
"entry": [
"src/index.ts"
],
"format": [
"cjs",
"esm"
],
"dts": true,
"splitting": false,
"sourcemap": true,
"clean": true,
"minify": false
}
}
+3129 -3482
View File
File diff suppressed because it is too large Load Diff