From 68e05f6cb1f2754cc152638f22f7c4fdb44f70bc Mon Sep 17 00:00:00 2001 From: Patrik Oldsberg Date: Tue, 21 May 2024 11:51:02 +0200 Subject: [PATCH] events-backend: connection event handlers + auth Signed-off-by: Patrik Oldsberg --- plugins/events-backend/dev/index.ts | 17 ++- .../src/service/EventsPlugin.ts | 5 +- .../src/service/hub/EventHub.ts | 129 +++++++++++++----- 3 files changed, 110 insertions(+), 41 deletions(-) diff --git a/plugins/events-backend/dev/index.ts b/plugins/events-backend/dev/index.ts index 95ed3e0be8..f5eb40ad50 100644 --- a/plugins/events-backend/dev/index.ts +++ b/plugins/events-backend/dev/index.ts @@ -59,9 +59,10 @@ backend.add( events: eventsServiceRef, logger: coreServices.logger, discovery: coreServices.discovery, + auth: coreServices.auth, rootLifecycle: coreServices.rootLifecycle, }, - async init({ events, logger, discovery, rootLifecycle }) { + async init({ events, logger, discovery, rootLifecycle, auth }) { events.subscribe({ id: 'test-1', topics: ['test'], @@ -74,7 +75,15 @@ backend.add( logger.info('Started!'); const baseUrl = await discovery.getBaseUrl('events'); console.log(`DEBUG: baseUrl=`, baseUrl); - const ws = new WebSocket(`${baseUrl}/hub/connect`); + const { token } = await auth.getPluginRequestToken({ + onBehalfOf: await auth.getOwnServiceCredentials(), + targetPluginId: 'events', + }); + const ws = new WebSocket(`${baseUrl}/hub/connect`, { + headers: { + Authorization: `Bearer ${token}`, + }, + }); ws.onopen = () => { console.log('DEBUG: ws.onopen'); ws.send('derp!'); @@ -82,8 +91,8 @@ backend.add( ws.onmessage = event => { console.log(`DEBUG: event=`, event.data); }; - ws.onerror = error => { - console.log(`Client error`, String(error)); + ws.onerror = event => { + console.log(`Client error`, event.error); }; }); }, diff --git a/plugins/events-backend/src/service/EventsPlugin.ts b/plugins/events-backend/src/service/EventsPlugin.ts index 0e115c9772..ecb234066d 100644 --- a/plugins/events-backend/src/service/EventsPlugin.ts +++ b/plugins/events-backend/src/service/EventsPlugin.ts @@ -76,9 +76,10 @@ export const eventsPlugin = createBackendPlugin({ config: coreServices.rootConfig, events: eventsServiceRef, logger: coreServices.logger, + httpAuth: coreServices.httpAuth, router: coreServices.httpRouter, }, - async init({ config, events, logger, router }) { + async init({ config, events, logger, httpAuth, router }) { const ingresses = Object.fromEntries( extensionPoint.httpPostIngresses.map(ingress => [ ingress.topic, @@ -95,7 +96,7 @@ export const eventsPlugin = createBackendPlugin({ const eventsRouter = Router(); http.bind(eventsRouter); - const hub = await EventHub.create({ logger }); + const hub = await EventHub.create({ logger, httpAuth }); eventsRouter.use('/hub', hub.handler()); router.use(eventsRouter); diff --git a/plugins/events-backend/src/service/hub/EventHub.ts b/plugins/events-backend/src/service/hub/EventHub.ts index d98c9e4d3d..cc99a2b2dd 100644 --- a/plugins/events-backend/src/service/hub/EventHub.ts +++ b/plugins/events-backend/src/service/hub/EventHub.ts @@ -14,11 +14,17 @@ * limitations under the License. */ -import { LoggerService } from '@backstage/backend-plugin-api'; +import { + BackstageCredentials, + BackstageServicePrincipal, + HttpAuthService, + LoggerService, +} from '@backstage/backend-plugin-api'; import { Handler } from 'express'; import Router from 'express-promise-router'; import { Socket } from 'net'; -import { WebSocketServer, type WebSocket } from 'ws'; +import { STATUS_CODES } from 'http'; +import { WebSocketServer, type WebSocket, RawData } from 'ws'; /** * Manages a single WebSocket connection. @@ -30,63 +36,90 @@ class EventClientConnection { ws: WebSocket; socket: Socket; logger: LoggerService; + credentials: BackstageCredentials; }) { - const { ws } = options; + const { ws, credentials } = options; const id = Math.random().toString(36).slice(2, 10); - const logger = options.logger.child({ connection: id }); - - ws.addListener('ping', () => { - ws.pong(); + const logger = options.logger.child({ + connection: id, + subject: credentials.principal.subject, }); - ws.onmessage = event => { - logger.debug(`Message from client: ${JSON.stringify(event.data)}`); - }; - ws.send('hello there!'); + const conn = new EventClientConnection(id, ws, logger, credentials); - const conn = new EventClientConnection(id, ws, options.socket, logger); + ws.addListener('close', conn.#handleClose); + ws.addListener('error', conn.#handleError); + ws.addListener('message', conn.#handleMessage); + ws.addListener('ping', () => ws.pong()); - logger.info(`New ${conn}`); + logger.info( + `New event client connection from '${options.socket.remoteAddress}'`, + ); return conn; } readonly #id: string; readonly #ws: WebSocket; - readonly #socket: Socket; readonly #logger: LoggerService; + readonly #credentials: BackstageCredentials; constructor( id: string, ws: WebSocket, - socket: Socket, logger: LoggerService, + credentials: BackstageCredentials, ) { this.#id = id; this.#ws = ws; - this.#socket = socket; 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 = (data: RawData, isBinary: boolean) => { + console.log(`DEBUG: isBinary=${isBinary} data=${data}`); + }; + close() { + this.#removeListeners(); this.#ws.close(); - this.#logger.info(`Closed ${this}`); + 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},addr=${ - this.#socket.remoteAddress + return `eventClientConnection{id=${this.#id},subject=${ + this.#credentials.principal.subject }}`; } } export class EventHub { - static async create(options: { logger: LoggerService }) { + static async create(options: { + logger: LoggerService; + httpAuth: HttpAuthService; + }) { + const { httpAuth } = options; const logger = options.logger.child({ type: 'EventHub' }); const router = Router(); @@ -94,11 +127,12 @@ export class EventHub { noServer: true, clientTracking: false, }); + server.on('error', error => { logger.error(`WebSocket server error`, error); }); - const hub = new EventHub(server, router, logger); + const hub = new EventHub(server, router, logger, httpAuth); router.get('/connect', hub.#handleGetConnect); @@ -108,6 +142,7 @@ export class EventHub { readonly #server: WebSocketServer; readonly #handler: Handler; readonly #logger: LoggerService; + readonly #httpAuth: HttpAuthService; #connections = new Map(); @@ -115,30 +150,54 @@ export class EventHub { server: WebSocketServer, handler: Handler, logger: LoggerService, + httpAuth: HttpAuthService, ) { this.#server = server; this.#handler = handler; this.#logger = logger; + this.#httpAuth = httpAuth; } handler(): Handler { return this.#handler; } - #handleGetConnect: Handler = (req, _res) => { - this.#server.handleUpgrade( - req, - req.socket, - Buffer.alloc(0), - (ws, { socket }) => { - const conn = EventClientConnection.create({ - ws, - socket, - logger: this.#logger, - }); - this.#connections.set(conn.id, conn); - }, - ); + #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, + }); + 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); + } }; close() {