From ee52e38a04aca350b2301a6da6a77babce728238 Mon Sep 17 00:00:00 2001 From: Patrik Oldsberg Date: Wed, 22 May 2024 18:34:56 +0200 Subject: [PATCH] events-backend: add hub http subscription API + listen by subscription ID Signed-off-by: Patrik Oldsberg --- plugins/events-backend/dev/index.ts | 49 ++++++++++++--- .../src/service/hub/EventHub.ts | 63 ++++++++++++++++--- 2 files changed, 96 insertions(+), 16 deletions(-) diff --git a/plugins/events-backend/dev/index.ts b/plugins/events-backend/dev/index.ts index 33bb01a985..77aa379744 100644 --- a/plugins/events-backend/dev/index.ts +++ b/plugins/events-backend/dev/index.ts @@ -79,6 +79,39 @@ backend.add( onBehalfOf: await auth.getOwnServiceCredentials(), targetPluginId: 'events', }); + + const subRes = await fetch(`${baseUrl}/hub/subscriptions/123`, { + method: 'PUT', + headers: { + Authorization: `Bearer ${token}`, + 'Content-Type': 'application/json', + }, + body: JSON.stringify({ + topics: ['test'], + }), + }); + console.log( + `DEBUG: sub create req = ${subRes.status} ${subRes.statusText}`, + ); + + const poll = async () => { + const res = await fetch(`${baseUrl}/hub/subscriptions/123`, { + headers: { + Authorization: `Bearer ${token}`, + 'Content-Type': 'application/json', + }, + }); + + const data = res.status === 200 && (await res.json()); + console.log( + `DEBUG: sub poll req = ${res.status} ${res.statusText}`, + data, + ); + poll(); + }; + + poll(); + const ws = new WebSocket(`${baseUrl}/hub/connect`, { headers: { Authorization: `Bearer ${token}`, @@ -86,14 +119,14 @@ backend.add( }); ws.onopen = () => { console.log('DEBUG: ws.onopen'); - ws.send( - JSON.stringify([ - 'req', - 1, - 'subscribe', - { id: 'derp', topics: ['test'] }, - ]), - ); + // ws.send( + // JSON.stringify([ + // 'req', + // 1, + // 'subscribe', + // { id: 'derp', topics: ['test'] }, + // ]), + // ); setTimeout(() => { console.log(`DEBUG: publish!`); ws.send( diff --git a/plugins/events-backend/src/service/hub/EventHub.ts b/plugins/events-backend/src/service/hub/EventHub.ts index 7cb68ce808..c4d0a785d3 100644 --- a/plugins/events-backend/src/service/hub/EventHub.ts +++ b/plugins/events-backend/src/service/hub/EventHub.ts @@ -20,7 +20,7 @@ import { HttpAuthService, LoggerService, } from '@backstage/backend-plugin-api'; -import { Handler } from 'express'; +import express, { Handler } from 'express'; import Router from 'express-promise-router'; import { Socket } from 'net'; import { STATUS_CODES } from 'http'; @@ -113,7 +113,7 @@ type EventHubStore = { readSubscription(id: string): Promise<{ events: EventParams[] }>; listen( - topicIds: string[], + subscriptionId: string, onNotify: (topicId: string) => void, ): Promise<() => void>; }; @@ -191,10 +191,14 @@ class MemoryEventHubStore implements EventHubStore { } async listen( - topicIds: string[], + subscriptionId: string, onNotify: (topicId: string) => void, ): Promise<() => void> { - const listener = { topics: new Set(topicIds), notify: onNotify }; + const sub = this.#subscribers.get(subscriptionId); + if (!sub) { + throw new Error(`Subscription not found`); + } + const listener = { topics: sub.topics, notify: onNotify }; this.#listeners.add(listener); return () => { this.#listeners.delete(listener); @@ -390,8 +394,15 @@ export class EventHub { const hub = new EventHub(server, router, logger, httpAuth); + // WS router.get('/connect', hub.#handleGetConnect); + router.use(express.json()); + + // Long-polling + router.get('/subscriptions/:id', hub.#handleGetSubscription); + router.put('/subscriptions/:id', hub.#handlePutSubscription); + return hub; } @@ -474,10 +485,7 @@ export class EventHub { }, ); - const removeListener = await this.#store.listen( - params.topics, - read, - ); + const removeListener = await this.#store.listen(params.id, read); ws.addListener('close', removeListener); await read(); @@ -518,6 +526,45 @@ export class EventHub { } }; + #handleGetSubscription: Handler = async (req, res) => { + const credentials = await this.#httpAuth.credentials(req, { + allow: ['service'], + }); + const id = req.params.id; + + const { events } = await this.#store.readSubscription(id); + + this.#logger.info( + `Reading subscription '${id}' resulted in ${events.length} events`, + { subject: credentials.principal.subject }, + ); + + if (events.length > 0) { + res.json({ events }); + return; + } + + this.#store.listen(id, () => { + res.status(204).end(); + }); + }; + + #handlePutSubscription: Handler = async (req, res) => { + const credentials = await this.#httpAuth.credentials(req, { + allow: ['service'], + }); + const id = req.params.id; + + await this.#store.upsertSubscription(id, req.body.topics); + + this.#logger.info( + `New subscription '${id}' topics='${req.body.topics.join("', '")}'`, + { subject: credentials.principal.subject }, + ); + + res.status(201).end(); + }; + close() { this.#connections.forEach(conn => conn.close()); this.#connections.clear();