diff --git a/plugins/events-backend/dev/index.ts b/plugins/events-backend/dev/index.ts index 47c2763a01..1c0d38215b 100644 --- a/plugins/events-backend/dev/index.ts +++ b/plugins/events-backend/dev/index.ts @@ -19,7 +19,6 @@ import { coreServices, createBackendPlugin, } from '@backstage/backend-plugin-api'; -import { WebSocket } from 'ws'; import { eventsServiceRef } from '@backstage/plugin-events-node'; import { DefaultApiClient } from '../../events-node/src/generated'; diff --git a/plugins/events-backend/package.json b/plugins/events-backend/package.json index aa788330ae..2f9db1b081 100644 --- a/plugins/events-backend/package.json +++ b/plugins/events-backend/package.json @@ -62,10 +62,7 @@ "@types/express": "^4.17.6", "express": "^4.17.1", "express-promise-router": "^4.1.0", - "winston": "^3.2.1", - "ws": "^8.17.0", - "zod": "^3.22.4", - "zod-validation-error": "^3.3.0" + "winston": "^3.2.1" }, "devDependencies": { "@backstage/backend-defaults": "workspace:^", diff --git a/plugins/events-backend/src/service/hub/EventHub.ts b/plugins/events-backend/src/service/hub/EventHub.ts index e8080ecb30..a014bb1add 100644 --- a/plugins/events-backend/src/service/hub/EventHub.ts +++ b/plugins/events-backend/src/service/hub/EventHub.ts @@ -14,96 +14,13 @@ * limitations under the License. */ -import { - BackstageCredentials, - BackstageServicePrincipal, - HttpAuthService, - LoggerService, -} from '@backstage/backend-plugin-api'; +import { HttpAuthService, LoggerService } from '@backstage/backend-plugin-api'; import { Handler } from 'express'; import Router from 'express-promise-router'; -import { Socket } from 'net'; -import { STATUS_CODES } from 'http'; -import { WebSocketServer, type WebSocket, RawData } from 'ws'; -import { z, ZodError } from 'zod'; -import { fromZodError } from 'zod-validation-error'; -import { JsonObject, JsonValue } from '@backstage/types'; -import { serializeError } from '@backstage/errors'; import { EventParams } from '@backstage/plugin-events-node'; import { spec, createOpenApiRouter } from '../../schema/openapi.generated'; import { internal } from '@backstage/backend-openapi-utils'; -/* - -# Protocol - -## Request/Response - -General request/response format used for all communication: - --> [type: 'req', id: number, method: string, params: JsonObject] -<- [type: 'res', id: number, status: 'resolved' | 'rejected', result: JsonObject] - -## Client -> Server - -### Subscribe - --> method: 'subscribe', params: { id: string, topics: string[] } -<- result: void - -### Publish - --> method: 'publish', params: { topic: string, payload: JsonObject } -<- result: void - -## Server -> Client - -### Event - --> method: 'event', params: { topic: string, payload: JsonObject } -<- result: void - -*/ - -const messageSchema = z.union([ - z.tuple([ - z.literal('req'), - z.number().int().gt(0), - z.string().min(1), - z.any(), - ]), - z.tuple([ - z.literal('res'), - z.number().int().gt(0), - z.enum(['resolved', 'rejected']), - z.any(), - ]), -]); -const subscribeParamsSchema = z.object({ - id: z.string().min(1), - topics: z.array(z.string().min(1)), -}); -const publishParamsSchema = z.object({ - topic: z.string().min(1), - payload: z.any(), -}); -const eventParamsSchema = z.object({ - events: z.array( - z.object({ - topic: z.string().min(1), - payload: z.any(), - metadata: z.any().optional(), - }), - ), -}); - -function errorToJson(error: Error): JsonObject { - if (error.name === 'ZodError') { - return serializeError(fromZodError(error as ZodError)); - } - return serializeError(error); -} - type EventHubStore = { publish(options: { params: EventParams; @@ -216,174 +133,6 @@ class MemoryEventHubStore implements EventHubStore { } } -/** - * Manages a single WebSocket connection. - * - * @internal - */ -class EventClientConnection { - static create(options: { - ws: WebSocket; - socket: Socket; - logger: LoggerService; - credentials: BackstageCredentials; - }) { - const { ws, credentials } = options; - - const id = Math.random().toString(36).slice(2, 10); - const logger = options.logger.child({ - connection: id, - subject: credentials.principal.subject, - }); - - const conn = new EventClientConnection(id, ws, logger, credentials); - - ws.addListener('close', conn.#handleClose); - ws.addListener('error', conn.#handleError); - ws.addListener('message', conn.#handleMessage); - - logger.info( - `New event client connection from '${options.socket.remoteAddress}'`, - ); - - return conn; - } - - readonly #id: string; - readonly #ws: WebSocket; - readonly #logger: LoggerService; - readonly #credentials: BackstageCredentials; - - #seq = 1; - readonly #pendingRequests = new Map< - number, - { resolve(result: unknown): void; reject(error: unknown): void } - >(); - readonly #requestHandlers = new Map< - string, - { - schema: z.ZodType; - handler: (params: any) => unknown; - } - >(); - - constructor( - id: string, - ws: WebSocket, - logger: LoggerService, - credentials: BackstageCredentials, - ) { - this.#id = id; - this.#ws = ws; - this.#logger = logger; - this.#credentials = credentials; - } - - get id() { - return this.#id; - } - - #handleClose = (code: number, reason: Buffer) => { - this.#removeListeners(); - this.#logger.info(`Remote closed code=${code} reason=${reason}`); - }; - - #handleError = (error: Error) => { - this.#removeListeners(); - this.#logger.error(`WebSocket error`, error); - }; - - #handleMessage = (rawData: RawData, isBinary: boolean) => { - if (isBinary) { - return; - } - try { - const data = Array.isArray(rawData) - ? Buffer.concat(rawData) - : Buffer.from(rawData); - const message = messageSchema.parse(JSON.parse(data.toString('utf8'))); - - if (message[0] === 'req') { - const [, seq, method, params] = message; - const handler = this.#requestHandlers.get(method); - if (!handler) { - throw new Error(`Unknown method '${method}'`); - } - - try { - const parsedParams = handler.schema.parse(params); - - Promise.resolve(handler.handler(parsedParams)).then( - result => { - this.#sendMessage('res', seq, 'resolved', result ?? null); - }, - error => { - this.#sendMessage('res', seq, 'rejected', errorToJson(error)); - }, - ); - } catch (error) { - this.#sendMessage('res', seq, false, errorToJson(error)); - } - } else if (message[0] === 'res') { - const [, seq, success, result] = message; - const pendingRequest = this.#pendingRequests.get(seq); - if (!pendingRequest) { - throw new Error(`Received response for unknown request seq=${seq}`); - } - this.#pendingRequests.delete(seq); - if (success) { - pendingRequest.resolve(result); - } else { - pendingRequest.reject(result); - } - } - } catch (error) { - this.#logger.error('Invalid message received', error); - } - }; - - addRequestHandler( - method: string, - schema: z.ZodType, - handler: (params: TParams) => unknown, - ) { - this.#requestHandlers.set(method, { schema, handler }); - } - - async request( - method: string, - params: TReq, - ): Promise { - return new Promise((resolve, reject) => { - const seq = this.#seq++; - this.#pendingRequests.set(seq, { resolve, reject }); - this.#sendMessage('req', seq, method, params); - }); - } - - #sendMessage(...message: JsonValue[]) { - this.#ws.send(JSON.stringify(message)); - } - - close() { - this.#removeListeners(); - this.#ws.close(); - this.#logger.info(`Closed connection`); - } - - #removeListeners() { - this.#ws.removeListener('close', this.#handleClose); - this.#ws.removeListener('error', this.#handleError); - this.#ws.removeListener('message', this.#handleMessage); - } - - toString() { - return `eventClientConnection{id=${this.#id},subject=${ - this.#credentials.principal.subject - }}`; - } -} - export class EventHub { static async create(options: { logger: LoggerService; @@ -393,19 +142,7 @@ export class EventHub { const logger = options.logger.child({ type: 'EventHub' }); const router = Router(); - const server = new WebSocketServer({ - noServer: true, - clientTracking: false, - }); - - server.on('error', error => { - logger.error(`WebSocket server error`, error); - }); - - const hub = new EventHub(server, router, logger, httpAuth); - - // WS - router.get('/hub/connect', hub.#handleGetConnect); + const hub = new EventHub(router, logger, httpAuth); const apiRouter = await createOpenApiRouter(); @@ -426,21 +163,16 @@ export class EventHub { return hub; } - readonly #server: WebSocketServer; readonly #handler: Handler; readonly #logger: LoggerService; readonly #httpAuth: HttpAuthService; readonly #store: EventHubStore; - #connections = new Map(); - private constructor( - server: WebSocketServer, handler: Handler, logger: LoggerService, httpAuth: HttpAuthService, ) { - this.#server = server; this.#handler = handler; this.#logger = logger; this.#httpAuth = httpAuth; @@ -451,101 +183,6 @@ export class EventHub { return this.#handler; } - #handleGetConnect: Handler = async (req, _res) => { - try { - const credentials = await this.#httpAuth.credentials(req, { - allow: ['service'], - }); - - this.#server.handleUpgrade( - req, - req.socket, - Buffer.alloc(0), - (ws, { socket }) => { - const conn = EventClientConnection.create({ - ws, - socket, - logger: this.#logger, - credentials, - }); - conn.addRequestHandler( - 'subscribe', - subscribeParamsSchema, - async params => { - await this.#store.upsertSubscription(params.id, params.topics); - - this.#logger.info( - `New subscription '${params.id}' topics='${params.topics.join( - "', '", - )}'`, - ); - - const read = () => - this.#store.readSubscription(params.id).then( - ({ events }) => { - if (events.length > 0) { - conn - .request, void>( - 'events', - { events }, - ) - .catch(error => { - this.#logger.error( - `Failed to send events to subscription ${params.id}`, - error, - ); - }); - } - }, - error => { - this.#logger.error( - `Failed to read subscription ${params.id}`, - error, - ); - }, - ); - - const removeListener = await this.#store.listen(params.id, read); - ws.addListener('close', removeListener); - - await read(); - }, - ); - conn.addRequestHandler( - 'publish', - publishParamsSchema, - async params => { - await this.#store.publish({ - params: { - topic: params.topic, - eventPayload: params.payload, - }, - subscriberIds: [], - }); - this.#logger.info(`Published event to '${params.topic}'`); - }, - ); - this.#connections.set(conn.id, conn); - }, - ); - } catch (error) { - let status = 500; - if (error.name === 'AuthenticationError') { - status = 401; - } else if (error.name === 'NotAllowedError') { - status = 403; - } - req.socket.write( - `HTTP/1.1 ${status} ${STATUS_CODES[status]}\r\n` + - 'Upgrade: WebSocket\r\n' + - 'Connection: Upgrade\r\n' + - '\r\n', - ); - req.socket.destroy(); - this.#logger.info('WebSocket upgrade failed', error); - } - }; - #handlePostEvents: internal.DocRequestHandler< typeof spec, '/hub/events', @@ -613,10 +250,4 @@ export class EventHub { res.status(201).end(); }; - - close() { - this.#connections.forEach(conn => conn.close()); - this.#connections.clear(); - this.#server.close(); - } } diff --git a/yarn.lock b/yarn.lock index 6636108958..2aeb54be7f 100644 --- a/yarn.lock +++ b/yarn.lock @@ -6154,9 +6154,6 @@ __metadata: express-promise-router: ^4.1.0 supertest: ^7.0.0 winston: ^3.2.1 - ws: ^8.17.0 - zod: ^3.22.4 - zod-validation-error: ^3.3.0 languageName: unknown linkType: soft @@ -44671,12 +44668,12 @@ __metadata: languageName: node linkType: hard -"zod-validation-error@npm:^3.0.3, zod-validation-error@npm:^3.3.0": - version: 3.3.0 - resolution: "zod-validation-error@npm:3.3.0" +"zod-validation-error@npm:^3.0.3": + version: 3.1.0 + resolution: "zod-validation-error@npm:3.1.0" peerDependencies: zod: ^3.18.0 - checksum: cbf81ecd27df675d72883b69833565af787302e70ad970ae4a5dab84e1cb8739cedf094b35f7f4b78307adaadb7cab0c0a8f7debeb6516e3fee998a3d4e13422 + checksum: 84df01c91d594701eaf7f5f007be881e47f7adef2e3f3765f7be031cb78033f9be0924273106cb81b586d8020da9885dbb81b3da363f00a51df00f26274f2b23 languageName: node linkType: hard