From 335443229f9ce802a3b1dbd2e9a73a718a1069e0 Mon Sep 17 00:00:00 2001 From: Patrik Oldsberg Date: Mon, 20 May 2024 18:49:44 +0200 Subject: [PATCH] events-backend: initial WebSocket server Signed-off-by: Patrik Oldsberg --- plugins/events-backend/dev/index.ts | 38 ++++++-- plugins/events-backend/package.json | 3 +- .../src/service/EventsPlugin.ts | 5 ++ .../src/service/hub/EventHub.ts | 86 +++++++++++++++++++ .../events-backend/src/service/hub/index.ts | 17 ++++ yarn.lock | 1 + 6 files changed, 140 insertions(+), 10 deletions(-) create mode 100644 plugins/events-backend/src/service/hub/EventHub.ts create mode 100644 plugins/events-backend/src/service/hub/index.ts diff --git a/plugins/events-backend/dev/index.ts b/plugins/events-backend/dev/index.ts index c0626543c8..95ed3e0be8 100644 --- a/plugins/events-backend/dev/index.ts +++ b/plugins/events-backend/dev/index.ts @@ -19,6 +19,7 @@ import { coreServices, createBackendPlugin, } from '@backstage/backend-plugin-api'; +import { WebSocket } from 'ws'; import { eventsServiceRef } from '@backstage/plugin-events-node'; const backend = createBackend(); @@ -35,14 +36,14 @@ backend.add( logger: coreServices.logger, }, async init({ events, logger }) { - setInterval(() => { - logger.info(`Publishing event to topic 'test'`); - events.publish({ - eventPayload: { foo: 'bar' }, - topic: 'test', - metadata: { meta: 'baz' }, - }); - }, 5000); + // setInterval(() => { + // logger.info(`Publishing event to topic 'test'`); + // events.publish({ + // eventPayload: { foo: 'bar' }, + // topic: 'test', + // metadata: { meta: 'baz' }, + // }); + // }, 5000); }, }); }, @@ -57,8 +58,10 @@ backend.add( deps: { events: eventsServiceRef, logger: coreServices.logger, + discovery: coreServices.discovery, + rootLifecycle: coreServices.rootLifecycle, }, - async init({ events, logger }) { + async init({ events, logger, discovery, rootLifecycle }) { events.subscribe({ id: 'test-1', topics: ['test'], @@ -66,6 +69,23 @@ backend.add( logger.info(`Received event: ${JSON.stringify(event, null, 2)}`); }, }); + + rootLifecycle.addStartupHook(async () => { + logger.info('Started!'); + const baseUrl = await discovery.getBaseUrl('events'); + console.log(`DEBUG: baseUrl=`, baseUrl); + const ws = new WebSocket(`${baseUrl}/hub/connect`); + ws.onopen = () => { + console.log('DEBUG: ws.onopen'); + ws.send('derp!'); + }; + ws.onmessage = event => { + console.log(`DEBUG: event=`, event.data); + }; + ws.onerror = error => { + console.log(`Client error`, String(error)); + }; + }); }, }); }, diff --git a/plugins/events-backend/package.json b/plugins/events-backend/package.json index c25109cf82..6665ff1e52 100644 --- a/plugins/events-backend/package.json +++ b/plugins/events-backend/package.json @@ -58,7 +58,8 @@ "@types/express": "^4.17.6", "express": "^4.17.1", "express-promise-router": "^4.1.0", - "winston": "^3.2.1" + "winston": "^3.2.1", + "ws": "^8.17.0" }, "devDependencies": { "@backstage/backend-defaults": "workspace:^", diff --git a/plugins/events-backend/src/service/EventsPlugin.ts b/plugins/events-backend/src/service/EventsPlugin.ts index d95dbf3efc..0e115c9772 100644 --- a/plugins/events-backend/src/service/EventsPlugin.ts +++ b/plugins/events-backend/src/service/EventsPlugin.ts @@ -28,6 +28,7 @@ import { } from '@backstage/plugin-events-node'; import Router from 'express-promise-router'; import { HttpPostIngressEventPublisher } from './http'; +import { EventHub } from './hub'; class EventsExtensionPointImpl implements EventsExtensionPoint { #httpPostIngresses: HttpPostIngressOptions[] = []; @@ -93,6 +94,10 @@ export const eventsPlugin = createBackendPlugin({ }); const eventsRouter = Router(); http.bind(eventsRouter); + + const hub = await EventHub.create({ logger }); + eventsRouter.use('/hub', hub.handler()); + router.use(eventsRouter); router.addAuthPolicy({ allow: 'unauthenticated', diff --git a/plugins/events-backend/src/service/hub/EventHub.ts b/plugins/events-backend/src/service/hub/EventHub.ts new file mode 100644 index 0000000000..b1ef75d9e6 --- /dev/null +++ b/plugins/events-backend/src/service/hub/EventHub.ts @@ -0,0 +1,86 @@ +/* + * Copyright 2024 The Backstage Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +import { LoggerService } from '@backstage/backend-plugin-api'; +import { Handler } from 'express'; +import Router from 'express-promise-router'; +import { WebSocketServer, type WebSocket } from 'ws'; + +export class EventHub { + static async create(options: { logger: LoggerService }) { + 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); + + router.get('/connect', hub.#handleGetConnect); + + return hub; + } + + readonly #server: WebSocketServer; + readonly #handler: Handler; + readonly #logger: LoggerService; + + #connections = new Set(); + + private constructor( + server: WebSocketServer, + handler: Handler, + logger: LoggerService, + ) { + this.#server = server; + this.#handler = handler; + this.#logger = logger; + } + + handler(): Handler { + return this.#handler; + } + + #handleGetConnect: Handler = (req, _res) => { + this.#server.handleUpgrade(req, req.socket, Buffer.alloc(0), conn => { + const id = Math.random().toString(36).slice(2, 10); + const logger = this.#logger.child({ connection: id }); + + logger.info(`New connection from '${req.socket.remoteAddress}'`); + this.#connections.add(conn); + + conn.onmessage = event => { + logger.debug(`Message from client: ${JSON.stringify(event.data)}`); + }; + conn.send('hello there!'); + + conn.addListener('ping', () => { + conn.pong(); + }); + }); + }; + + close() { + this.#connections.forEach(conn => conn.close()); + this.#connections.clear(); + this.#server.close(); + } +} diff --git a/plugins/events-backend/src/service/hub/index.ts b/plugins/events-backend/src/service/hub/index.ts new file mode 100644 index 0000000000..4acc848298 --- /dev/null +++ b/plugins/events-backend/src/service/hub/index.ts @@ -0,0 +1,17 @@ +/* + * Copyright 2024 The Backstage Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +export { EventHub } from './EventHub'; diff --git a/yarn.lock b/yarn.lock index cee8a0c9ac..18e926a14e 100644 --- a/yarn.lock +++ b/yarn.lock @@ -6150,6 +6150,7 @@ __metadata: express-promise-router: ^4.1.0 supertest: ^7.0.0 winston: ^3.2.1 + ws: ^8.17.0 languageName: unknown linkType: soft