events-backend: add hub http subscription API + listen by subscription ID

Signed-off-by: Patrik Oldsberg <poldsberg@gmail.com>
This commit is contained in:
Patrik Oldsberg
2024-05-22 18:34:56 +02:00
parent a90ce4aa39
commit ee52e38a04
2 changed files with 96 additions and 16 deletions
+41 -8
View File
@@ -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(
@@ -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();