Merge pull request #16678 from RoadieHQ/add-export-for-in-memory-broker
Add export for in memory broker
This commit is contained in:
@@ -6,7 +6,7 @@ This plugin provides the wiring of all extension points
|
||||
for managing events as defined by [plugin-events-node](../events-node)
|
||||
including backend plugin `EventsPlugin` and `EventsBackend`.
|
||||
|
||||
Additionally, it uses a simple in-memory implementation for
|
||||
Additionally, it uses a simple in-process implementation for
|
||||
the `EventBroker` by default which you can replace with a more sophisticated
|
||||
implementation of your choice as you need (e.g., via module).
|
||||
|
||||
@@ -24,10 +24,89 @@ to the used event broker.
|
||||
yarn add --cwd packages/backend @backstage/plugin-events-backend
|
||||
```
|
||||
|
||||
Add a file [`packages/backend/src/plugins/events.ts`](../../packages/backend/src/plugins/events.ts)
|
||||
to your Backstage project.
|
||||
### Event Broker
|
||||
|
||||
There, you can add all publishers, subscribers, etc. you want.
|
||||
First you will need to add and implementation of the `EventBroker` interface to the backend plugin environment.
|
||||
This will allow event broker instance any backend plugins to publish and subscribe to events in order to communicate
|
||||
between them.
|
||||
|
||||
Add the following to `makeCreateEnv`
|
||||
|
||||
```diff
|
||||
// packages/backend/src/index.ts
|
||||
+ const eventBroker = new DefaultEventBroker(root.child({ type: 'plugin' }));
|
||||
```
|
||||
|
||||
Then update plugin environment to include the event broker.
|
||||
|
||||
```diff
|
||||
// packages/backend/src/types.ts
|
||||
+ eventBroker: EventBroker;
|
||||
```
|
||||
|
||||
### Publishing and Subscribing to events with the broker
|
||||
|
||||
Backend plugins are passed the event broker in the plugin environment at startup of the application. The plugin can
|
||||
make use of this to communicate between parts of the application.
|
||||
|
||||
Here is an example of a plugin publishing a payload to a topic.
|
||||
|
||||
```typescript jsx
|
||||
export default async function createPlugin(
|
||||
env: PluginEnvironment,
|
||||
): Promise<Router> {
|
||||
env.eventBroker.publish({
|
||||
topic: 'publish.example',
|
||||
eventPayload: { message: 'Hello, World!' },
|
||||
metadata: {},
|
||||
});
|
||||
}
|
||||
```
|
||||
|
||||
Here is an example of a plugin subscribing to a topic.
|
||||
|
||||
```typescript jsx
|
||||
export default async function createPlugin(
|
||||
env: PluginEnvironment,
|
||||
): Promise<Router> {
|
||||
env.eventBroker.subscribe([
|
||||
{
|
||||
supportsEventTopics: ['publish.example'],
|
||||
onEvent: async (params: EventParams) => {
|
||||
env.logger.info(`receieved ${params.topic} event`);
|
||||
},
|
||||
},
|
||||
]);
|
||||
}
|
||||
```
|
||||
|
||||
### Implementing an `EventSubscriber` class
|
||||
|
||||
More complex solutions might need the creation of a class that implements the `EventSubscriber` interface. e.g.
|
||||
|
||||
```typescript jsx
|
||||
import { EventSubscriber } from "./EventSubscriber";
|
||||
|
||||
class ExampleSubscriber implements EventSubscriber {
|
||||
...
|
||||
|
||||
supportsEventTopics() {
|
||||
return ['publish.example']
|
||||
}
|
||||
|
||||
async onEvent(params: EventParams) {
|
||||
env.logger.info(`receieved ${params.topic} event`)
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
### Events Backend
|
||||
|
||||
The events backend plugin provides a router to handler http events and publish the http requests onto the event
|
||||
broker.
|
||||
|
||||
To configure it add a file [`packages/backend/src/plugins/events.ts`](../../packages/backend/src/plugins/events.ts)
|
||||
to your Backstage project.
|
||||
|
||||
Additionally, add the events plugin to your backend.
|
||||
|
||||
@@ -38,54 +117,11 @@ Additionally, add the events plugin to your backend.
|
||||
// [...]
|
||||
+ const eventsEnv = useHotMemoize(module, () => createEnv('events'));
|
||||
// [...]
|
||||
+ apiRouter.use('/events', await events(eventsEnv, []));
|
||||
+ apiRouter.use('/events', await events(eventsEnv));
|
||||
// [...]
|
||||
```
|
||||
|
||||
### With Event-based Entity Providers
|
||||
|
||||
In case you use event-based `EntityProviders`,
|
||||
you may need something like the following:
|
||||
|
||||
```diff
|
||||
// packages/backend/src/index.ts
|
||||
- apiRouter.use('/events', await events(eventsEnv, []));
|
||||
+ apiRouter.use('/events', await events(eventsEnv, eventBasedEntityProviders));
|
||||
```
|
||||
|
||||
as well as a file
|
||||
[`packages/backend/src/plugins/catalogEventBasedProviders.ts`](../../packages/backend/src/plugins/catalogEventBasedProviders.ts)
|
||||
which contains event-based entity providers.
|
||||
|
||||
In case you don't have this dependency added yet:
|
||||
|
||||
```bash
|
||||
# From your Backstage root directory
|
||||
yarn add --cwd packages/backend @backstage/plugin-events-backend
|
||||
```
|
||||
|
||||
```diff
|
||||
// packages/backend/src/plugins/catalog.ts
|
||||
import { CatalogBuilder } from '@backstage/plugin-catalog-backend';
|
||||
+import { EntityProvider } from '@backstage/plugin-catalog-node';
|
||||
import { ScaffolderEntitiesProcessor } from '@backstage/plugin-scaffolder-backend';
|
||||
import { Router } from 'express';
|
||||
import { PluginEnvironment } from '../types';
|
||||
|
||||
export default async function createPlugin(
|
||||
env: PluginEnvironment,
|
||||
+ providers?: Array<EntityProvider>,
|
||||
): Promise<Router> {
|
||||
const builder = await CatalogBuilder.create(env);
|
||||
builder.addProcessor(new ScaffolderEntitiesProcessor());
|
||||
+ builder.addEntityProvider(providers ?? []);
|
||||
const { processingEngine, router } = await builder.build();
|
||||
await processingEngine.start();
|
||||
return router;
|
||||
}
|
||||
```
|
||||
|
||||
## Configuration
|
||||
#### Configuration
|
||||
|
||||
In order to create HTTP endpoints to receive events for a certain
|
||||
topic, you need to add them at your configuration:
|
||||
@@ -115,6 +151,34 @@ in combination with suitable event subscribers.
|
||||
|
||||
However, it is not limited to these use cases.
|
||||
|
||||
### Event-based Entity Providers
|
||||
|
||||
You can implement the `EventSubscriber` interface on an `EntityProviders` to allow it to handle events from other plugins e.g. the event backend plugin
|
||||
mentioned above.
|
||||
|
||||
Assuming you have configured the `eventBroker` into the `PluginEnvironment` you can pass the broker to the entity provider for it to subscribe.
|
||||
|
||||
```diff
|
||||
// packages/backend/src/plugins/catalog.ts
|
||||
import { CatalogBuilder } from '@backstage/plugin-catalog-backend';
|
||||
+import { DemoEventBasedEntityProvider } from './DemoEventBasedEntityProvider';
|
||||
import { ScaffolderEntitiesProcessor } from '@backstage/plugin-scaffolder-backend';
|
||||
import { Router } from 'express';
|
||||
import { PluginEnvironment } from '../types';
|
||||
|
||||
export default async function createPlugin(
|
||||
env: PluginEnvironment,
|
||||
): Promise<Router> {
|
||||
const builder = await CatalogBuilder.create(env);
|
||||
builder.addProcessor(new ScaffolderEntitiesProcessor());
|
||||
+ const demoProvider = new DemoEventBasedEntityProvider({ logger: env.logger, topics: ['example'], eventBroker: env.eventBroker });
|
||||
+ builder.addEntityProvider(demoProvider);
|
||||
const { processingEngine, router } = await builder.build();
|
||||
await processingEngine.start();
|
||||
return router;
|
||||
}
|
||||
```
|
||||
|
||||
## Use Cases
|
||||
|
||||
### Custom Event Broker
|
||||
|
||||
@@ -5,12 +5,24 @@
|
||||
```ts
|
||||
import { Config } from '@backstage/config';
|
||||
import { EventBroker } from '@backstage/plugin-events-node';
|
||||
import { EventParams } from '@backstage/plugin-events-node';
|
||||
import { EventPublisher } from '@backstage/plugin-events-node';
|
||||
import { EventSubscriber } from '@backstage/plugin-events-node';
|
||||
import express from 'express';
|
||||
import { HttpPostIngressOptions } from '@backstage/plugin-events-node';
|
||||
import { Logger } from 'winston';
|
||||
|
||||
// @public
|
||||
export class DefaultEventBroker implements EventBroker {
|
||||
constructor(logger: Logger);
|
||||
// (undocumented)
|
||||
publish(params: EventParams): Promise<void>;
|
||||
// (undocumented)
|
||||
subscribe(
|
||||
...subscribers: Array<EventSubscriber | Array<EventSubscriber>>
|
||||
): void;
|
||||
}
|
||||
|
||||
// @public
|
||||
export class EventsBackend {
|
||||
constructor(logger: Logger);
|
||||
|
||||
@@ -22,3 +22,4 @@
|
||||
|
||||
export { EventsBackend } from './service/EventsBackend';
|
||||
export { HttpPostIngressEventPublisher } from './service/http';
|
||||
export { DefaultEventBroker } from './service/DefaultEventBroker';
|
||||
|
||||
+4
-4
@@ -17,15 +17,15 @@
|
||||
import { getVoidLogger } from '@backstage/backend-common';
|
||||
import { TestEventSubscriber } from '@backstage/plugin-events-backend-test-utils';
|
||||
import { EventParams, EventSubscriber } from '@backstage/plugin-events-node';
|
||||
import { InMemoryEventBroker } from './InMemoryEventBroker';
|
||||
import { DefaultEventBroker } from './DefaultEventBroker';
|
||||
|
||||
const logger = getVoidLogger();
|
||||
|
||||
describe('InMemoryEventBroker', () => {
|
||||
describe('DefaultEventBroker', () => {
|
||||
it('passes events to interested subscribers', () => {
|
||||
const subscriber1 = new TestEventSubscriber('test1', ['topicA', 'topicB']);
|
||||
const subscriber2 = new TestEventSubscriber('test2', ['topicB', 'topicC']);
|
||||
const eventBroker = new InMemoryEventBroker(logger);
|
||||
const eventBroker = new DefaultEventBroker(logger);
|
||||
|
||||
eventBroker.subscribe(subscriber1);
|
||||
eventBroker.subscribe(subscriber2);
|
||||
@@ -86,7 +86,7 @@ describe('InMemoryEventBroker', () => {
|
||||
})();
|
||||
|
||||
const errorSpy = jest.spyOn(logger, 'error');
|
||||
const eventBroker = new InMemoryEventBroker(logger);
|
||||
const eventBroker = new DefaultEventBroker(logger);
|
||||
|
||||
eventBroker.subscribe(subscriber1);
|
||||
await eventBroker.publish({ topic, eventPayload: '1' });
|
||||
+4
-2
@@ -22,12 +22,14 @@ import {
|
||||
import { Logger } from 'winston';
|
||||
|
||||
/**
|
||||
* In-memory event broker which will pass the event to all registered subscribers
|
||||
* In process event broker which will pass the event to all registered subscribers
|
||||
* interested in it.
|
||||
* Events will not be persisted in any form.
|
||||
*
|
||||
* @public
|
||||
*/
|
||||
// TODO(pjungermann): add prom metrics? (see plugins/catalog-backend/src/util/metrics.ts, etc.)
|
||||
export class InMemoryEventBroker implements EventBroker {
|
||||
export class DefaultEventBroker implements EventBroker {
|
||||
constructor(private readonly logger: Logger) {}
|
||||
|
||||
private readonly subscribers: {
|
||||
@@ -20,7 +20,7 @@ import {
|
||||
EventSubscriber,
|
||||
} from '@backstage/plugin-events-node';
|
||||
import { Logger } from 'winston';
|
||||
import { InMemoryEventBroker } from './InMemoryEventBroker';
|
||||
import { DefaultEventBroker } from './DefaultEventBroker';
|
||||
|
||||
/**
|
||||
* A builder that helps wire up all component parts of the event management.
|
||||
@@ -33,7 +33,7 @@ export class EventsBackend {
|
||||
private subscribers: EventSubscriber[] = [];
|
||||
|
||||
constructor(logger: Logger) {
|
||||
this.eventBroker = new InMemoryEventBroker(logger);
|
||||
this.eventBroker = new DefaultEventBroker(logger);
|
||||
}
|
||||
|
||||
setEventBroker(eventBroker: EventBroker): EventsBackend {
|
||||
|
||||
@@ -29,7 +29,7 @@ import {
|
||||
EventSubscriber,
|
||||
HttpPostIngressOptions,
|
||||
} from '@backstage/plugin-events-node';
|
||||
import { InMemoryEventBroker } from './InMemoryEventBroker';
|
||||
import { DefaultEventBroker } from './DefaultEventBroker';
|
||||
import Router from 'express-promise-router';
|
||||
import { HttpPostIngressEventPublisher } from './http';
|
||||
|
||||
@@ -113,7 +113,7 @@ export const eventsPlugin = createBackendPlugin({
|
||||
router.use(eventsRouter);
|
||||
|
||||
const eventBroker =
|
||||
extensionPoint.eventBroker ?? new InMemoryEventBroker(winstonLogger);
|
||||
extensionPoint.eventBroker ?? new DefaultEventBroker(winstonLogger);
|
||||
|
||||
eventBroker.subscribe(extensionPoint.subscribers);
|
||||
[extensionPoint.publishers, http]
|
||||
|
||||
Reference in New Issue
Block a user