feat(events)!: migrate HttpPostIngressEventPublisher and eventsPlugin to use EventsService

Signed-off-by: Patrick Jungermann <Patrick.Jungermann@gmail.com>
This commit is contained in:
Patrick Jungermann
2024-01-23 20:28:35 +01:00
parent eff3ca9ddd
commit c4bd79422a
9 changed files with 132 additions and 113 deletions
+3 -4
View File
@@ -43,18 +43,17 @@ export class EventsBackend {
}
// @public
export class HttpPostIngressEventPublisher implements EventPublisher {
export class HttpPostIngressEventPublisher {
// (undocumented)
bind(router: express.Router): void;
// (undocumented)
static fromConfig(env: {
config: Config;
events: EventsService;
ingresses?: {
[topic: string]: Omit<HttpPostIngressOptions, 'topic'>;
};
logger: Logger;
logger: LoggerService;
}): HttpPostIngressEventPublisher;
// (undocumented)
setEventBroker(eventBroker: EventBroker): Promise<void>;
}
```
@@ -14,22 +14,27 @@
* limitations under the License.
*/
import { createBackendModule } from '@backstage/backend-plugin-api';
import { mockServices, startTestBackend } from '@backstage/backend-test-utils';
import {
TestEventBroker,
TestEventPublisher,
TestEventSubscriber,
} from '@backstage/plugin-events-backend-test-utils';
createBackendModule,
createServiceFactory,
} from '@backstage/backend-plugin-api';
import { mockServices, startTestBackend } from '@backstage/backend-test-utils';
import { eventsServiceRef } from '@backstage/plugin-events-node';
import { eventsExtensionPoint } from '@backstage/plugin-events-node/alpha';
import { TestEventsService } from '@backstage/plugin-events-backend-test-utils';
import request from 'supertest';
import { eventsPlugin } from './EventsPlugin';
describe('eventPlugin', () => {
describe('eventsPlugin', () => {
it('should be initialized properly', async () => {
const eventBroker = new TestEventBroker();
const publisher = new TestEventPublisher();
const subscriber = new TestEventSubscriber('sub', ['fake']);
const eventsService = new TestEventsService();
const eventsServiceFactory = createServiceFactory({
service: eventsServiceRef,
deps: {},
async factory({}) {
return eventsService;
},
});
const testModule = createBackendModule({
pluginId: 'events',
@@ -40,9 +45,9 @@ describe('eventPlugin', () => {
events: eventsExtensionPoint,
},
async init({ events }) {
events.setEventBroker(eventBroker);
events.addPublishers(publisher);
events.addSubscribers(subscriber);
events.addHttpPostIngress({
topic: 'fake-ext',
});
},
});
},
@@ -51,6 +56,7 @@ describe('eventPlugin', () => {
const { server } = await startTestBackend({
extensionPoints: [],
features: [
eventsServiceFactory(),
eventsPlugin(),
testModule(),
mockServices.logger.factory(),
@@ -66,18 +72,24 @@ describe('eventPlugin', () => {
],
});
expect(publisher.eventBroker).toBe(eventBroker);
expect(eventBroker.subscribed.length).toEqual(1);
expect(eventBroker.subscribed[0]).toBe(subscriber);
const response = await request(server)
const response1 = await request(server)
.post('/api/events/http/fake')
.timeout(1000)
.send({ test: 'fake' });
expect(response.status).toBe(202);
expect(response1.status).toBe(202);
expect(eventBroker.published.length).toEqual(1);
expect(eventBroker.published[0].topic).toEqual('fake');
expect(eventBroker.published[0].eventPayload).toEqual({ test: 'fake' });
const response2 = await request(server)
.post('/api/events/http/fake-ext')
.timeout(1000)
.send({ test: 'fake-ext' });
expect(response2.status).toBe(202);
expect(eventsService.published).toHaveLength(2);
expect(eventsService.published[0].topic).toEqual('fake');
expect(eventsService.published[0].eventPayload).toEqual({ test: 'fake' });
expect(eventsService.published[1].topic).toEqual('fake-ext');
expect(eventsService.published[1].eventPayload).toEqual({
test: 'fake-ext',
});
});
});
@@ -18,59 +18,42 @@ import {
createBackendPlugin,
coreServices,
} from '@backstage/backend-plugin-api';
import { loggerToWinstonLogger } from '@backstage/backend-common';
import {
eventsExtensionPoint,
EventsExtensionPoint,
} from '@backstage/plugin-events-node/alpha';
import {
EventBroker,
EventPublisher,
EventSubscriber,
eventsServiceRef,
HttpPostIngressOptions,
} from '@backstage/plugin-events-node';
import { DefaultEventBroker } from './DefaultEventBroker';
import Router from 'express-promise-router';
import { HttpPostIngressEventPublisher } from './http';
class EventsExtensionPointImpl implements EventsExtensionPoint {
#eventBroker: EventBroker | undefined;
#httpPostIngresses: HttpPostIngressOptions[] = [];
#publishers: EventPublisher[] = [];
#subscribers: EventSubscriber[] = [];
setEventBroker(eventBroker: EventBroker): void {
this.#eventBroker = eventBroker;
setEventBroker(_: any): void {
throw new Error(
'setEventBroker is not supported anymore; use eventsServiceRef instead',
);
}
addPublishers(
...publishers: Array<EventPublisher | Array<EventPublisher>>
): void {
this.#publishers.push(...publishers.flat());
addPublishers(_: any): void {
throw new Error(
'addPublishers is not supported anymore; use EventsService instead',
);
}
addSubscribers(
...subscribers: Array<EventSubscriber | Array<EventSubscriber>>
): void {
this.#subscribers.push(...subscribers.flat());
addSubscribers(_: any): void {
throw new Error(
'addSubscribers is not supported anymore; use EventsService instead',
);
}
addHttpPostIngress(options: HttpPostIngressOptions) {
this.#httpPostIngresses.push(options);
}
get eventBroker() {
return this.#eventBroker;
}
get publishers() {
return this.#publishers;
}
get subscribers() {
return this.#subscribers;
}
get httpPostIngresses() {
return this.#httpPostIngresses;
}
@@ -90,12 +73,11 @@ export const eventsPlugin = createBackendPlugin({
env.registerInit({
deps: {
config: coreServices.rootConfig,
events: eventsServiceRef,
logger: coreServices.logger,
router: coreServices.httpRouter,
},
async init({ config, logger, router }) {
const winstonLogger = loggerToWinstonLogger(logger);
async init({ config, events, logger, router }) {
const ingresses = Object.fromEntries(
extensionPoint.httpPostIngresses.map(ingress => [
ingress.topic,
@@ -105,20 +87,13 @@ export const eventsPlugin = createBackendPlugin({
const http = HttpPostIngressEventPublisher.fromConfig({
config,
events,
ingresses,
logger: winstonLogger,
logger,
});
const eventsRouter = Router();
http.bind(eventsRouter);
router.use(eventsRouter);
const eventBroker =
extensionPoint.eventBroker ?? new DefaultEventBroker(winstonLogger);
eventBroker.subscribe(extensionPoint.subscribers);
[extensionPoint.publishers, http]
.flat()
.forEach(publisher => publisher.setEventBroker(eventBroker));
},
});
},
@@ -16,7 +16,7 @@
import { getVoidLogger } from '@backstage/backend-common';
import { ConfigReader } from '@backstage/config';
import { TestEventBroker } from '@backstage/plugin-events-backend-test-utils';
import { TestEventsService } from '@backstage/plugin-events-backend-test-utils';
import express from 'express';
import Router from 'express-promise-router';
import request from 'supertest';
@@ -36,9 +36,11 @@ describe('HttpPostIngressEventPublisher', () => {
const router = Router();
const app = express().use(router);
const events = new TestEventsService();
const publisher = HttpPostIngressEventPublisher.fromConfig({
config,
events,
ingresses: {
testB: {},
},
@@ -46,9 +48,6 @@ describe('HttpPostIngressEventPublisher', () => {
});
publisher.bind(router);
const eventBroker = new TestEventBroker();
await publisher.setEventBroker(eventBroker);
const notFoundResponse = await request(app)
.post('/http/unknown')
.timeout(1000)
@@ -69,18 +68,18 @@ describe('HttpPostIngressEventPublisher', () => {
.send({ testB: 'data' });
expect(response2.status).toBe(202);
expect(eventBroker.published.length).toEqual(2);
expect(eventBroker.published[0].topic).toEqual('testA');
expect(eventBroker.published[0].eventPayload).toEqual({ testA: 'data' });
expect(eventBroker.published[0].metadata).toEqual(
expect(events.published).toHaveLength(2);
expect(events.published[0].topic).toEqual('testA');
expect(events.published[0].eventPayload).toEqual({ testA: 'data' });
expect(events.published[0].metadata).toEqual(
expect.objectContaining({
'content-type': 'application/json',
'x-custom-header': 'test-value',
}),
);
expect(eventBroker.published[1].topic).toEqual('testB');
expect(eventBroker.published[1].eventPayload).toEqual({ testB: 'data' });
expect(eventBroker.published[1].metadata).toEqual(
expect(events.published[1].topic).toEqual('testB');
expect(events.published[1].eventPayload).toEqual({ testB: 'data' });
expect(events.published[1].metadata).toEqual(
expect.objectContaining({
'content-type': 'application/json',
'x-custom-header': 'test-value',
@@ -99,9 +98,11 @@ describe('HttpPostIngressEventPublisher', () => {
const router = Router();
const app = express().use(router);
const events = new TestEventsService();
const publisher = HttpPostIngressEventPublisher.fromConfig({
config,
events,
ingresses: {
testB: {
validator: async (req, context) => {
@@ -146,9 +147,6 @@ describe('HttpPostIngressEventPublisher', () => {
});
publisher.bind(router);
const eventBroker = new TestEventBroker();
await publisher.setEventBroker(eventBroker);
const response1 = await request(app)
.post('/http/testA')
.timeout(1000)
@@ -191,12 +189,12 @@ describe('HttpPostIngressEventPublisher', () => {
expect(response6.status).toBe(403);
expect(response6.body).toEqual({});
expect(eventBroker.published.length).toEqual(2);
expect(eventBroker.published[0].topic).toEqual('testA');
expect(eventBroker.published[0].eventPayload).toEqual({ test: 'data' });
expect(eventBroker.published[1].topic).toEqual('testB');
expect(eventBroker.published[1].eventPayload).toEqual({ test: 'data' });
expect(eventBroker.published[1].metadata).toEqual(
expect(events.published).toHaveLength(2);
expect(events.published[0].topic).toEqual('testA');
expect(events.published[0].eventPayload).toEqual({ test: 'data' });
expect(events.published[1].topic).toEqual('testB');
expect(events.published[1].eventPayload).toEqual({ test: 'data' });
expect(events.published[1].metadata).toEqual(
expect.objectContaining({
'x-test-signature': 'testB-signature',
}),
@@ -205,10 +203,12 @@ describe('HttpPostIngressEventPublisher', () => {
it('without configuration', async () => {
const config = new ConfigReader({});
const events = new TestEventsService();
expect(() =>
HttpPostIngressEventPublisher.fromConfig({
config,
events,
logger,
}),
).not.toThrow();
@@ -15,16 +15,15 @@
*/
import { errorHandler } from '@backstage/backend-common';
import { LoggerService } from '@backstage/backend-plugin-api';
import { Config } from '@backstage/config';
import {
EventBroker,
EventPublisher,
EventsService,
HttpPostIngressOptions,
RequestValidator,
} from '@backstage/plugin-events-node';
import express from 'express';
import Router from 'express-promise-router';
import { Logger } from 'winston';
import { RequestValidationContextImpl } from './validation';
/**
@@ -34,13 +33,12 @@ import { RequestValidationContextImpl } from './validation';
* @public
*/
// TODO(pjungermann): add prom metrics? (see plugins/catalog-backend/src/util/metrics.ts, etc.)
export class HttpPostIngressEventPublisher implements EventPublisher {
private eventBroker?: EventBroker;
export class HttpPostIngressEventPublisher {
static fromConfig(env: {
config: Config;
events: EventsService;
ingresses?: { [topic: string]: Omit<HttpPostIngressOptions, 'topic'> };
logger: Logger;
logger: LoggerService;
}): HttpPostIngressEventPublisher {
const topics =
env.config.getOptionalStringArray('events.http.topics') ?? [];
@@ -54,11 +52,12 @@ export class HttpPostIngressEventPublisher implements EventPublisher {
}
});
return new HttpPostIngressEventPublisher(env.logger, ingresses);
return new HttpPostIngressEventPublisher(env.events, env.logger, ingresses);
}
private constructor(
private readonly logger: Logger,
private readonly events: EventsService,
private readonly logger: LoggerService,
private readonly ingresses: {
[topic: string]: Omit<HttpPostIngressOptions, 'topic'>;
},
@@ -68,10 +67,6 @@ export class HttpPostIngressEventPublisher implements EventPublisher {
router.use('/http', this.createRouter(this.ingresses));
}
async setEventBroker(eventBroker: EventBroker): Promise<void> {
this.eventBroker = eventBroker;
}
private createRouter(ingresses: {
[topic: string]: Omit<HttpPostIngressOptions, 'topic'>;
}): express.Router {
@@ -108,7 +103,7 @@ export class HttpPostIngressEventPublisher implements EventPublisher {
}
const eventPayload = request.body;
await this.eventBroker!.publish({
await this.events.publish({
topic,
eventPayload,
metadata: request.headers,