events-backend: initial WebSocket server

Signed-off-by: Patrik Oldsberg <poldsberg@gmail.com>
This commit is contained in:
Patrik Oldsberg
2024-05-20 18:49:44 +02:00
parent 85abf247ef
commit 335443229f
6 changed files with 140 additions and 10 deletions
+29 -9
View File
@@ -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));
};
});
},
});
},
+2 -1
View File
@@ -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:^",
@@ -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',
@@ -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<WebSocket>();
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();
}
}
@@ -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';
+1
View File
@@ -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