feat(events)!: migrate AwsSqsConsumingEventPublisher and its backend module to use EventsService
Signed-off-by: Patrick Jungermann <Patrick.Jungermann@gmail.com>
This commit is contained in:
@@ -0,0 +1,32 @@
|
||||
---
|
||||
'@backstage/plugin-events-backend-module-aws-sqs': minor
|
||||
---
|
||||
|
||||
BREAKING CHANGE: Migrate `AwsSqsConsumingEventPublisher` and its backend module to use `EventsService`.
|
||||
|
||||
Uses the `EventsService` instead of `EventBroker` at `AwsSqsConsumingEventPublisher`,
|
||||
dropping the use of `EventPublisher` including `setEventBroker(..)`.
|
||||
|
||||
Now, `AwsSqsConsumingEventPublisher.fromConfig` requires `events: EventsService` as option.
|
||||
|
||||
```diff
|
||||
const sqs = AwsSqsConsumingEventPublisher.fromConfig({
|
||||
config: env.config,
|
||||
+ events: env.events,
|
||||
logger: env.logger,
|
||||
scheduler: env.scheduler,
|
||||
});
|
||||
+ await Promise.all(sqs.map(publisher => publisher.start()));
|
||||
|
||||
// e.g. at packages/backend/src/plugins/events.ts
|
||||
- await new EventsBackend(env.logger)
|
||||
- .setEventBroker(env.eventBroker)
|
||||
- .addPublishers(sqs)
|
||||
- .start();
|
||||
|
||||
// or for other kinds of setups
|
||||
- await Promise.all(sqs.map(publisher => publisher.setEventBroker(eventBroker)));
|
||||
```
|
||||
|
||||
`eventsModuleAwsSqsConsumingEventPublisher` uses the `eventsServiceRef` as dependency,
|
||||
instead of `eventsExtensionPoint`.
|
||||
@@ -4,20 +4,20 @@
|
||||
|
||||
```ts
|
||||
import { Config } from '@backstage/config';
|
||||
import { EventBroker } from '@backstage/plugin-events-node';
|
||||
import { EventPublisher } from '@backstage/plugin-events-node';
|
||||
import { Logger } from 'winston';
|
||||
import { EventsService } from '@backstage/plugin-events-node';
|
||||
import { LoggerService } from '@backstage/backend-plugin-api';
|
||||
import { PluginTaskScheduler } from '@backstage/backend-tasks';
|
||||
|
||||
// @public
|
||||
export class AwsSqsConsumingEventPublisher implements EventPublisher {
|
||||
export class AwsSqsConsumingEventPublisher {
|
||||
// (undocumented)
|
||||
static fromConfig(env: {
|
||||
config: Config;
|
||||
logger: Logger;
|
||||
events: EventsService;
|
||||
logger: LoggerService;
|
||||
scheduler: PluginTaskScheduler;
|
||||
}): AwsSqsConsumingEventPublisher[];
|
||||
// (undocumented)
|
||||
setEventBroker(eventBroker: EventBroker): Promise<void>;
|
||||
start(): Promise<void>;
|
||||
}
|
||||
```
|
||||
|
||||
@@ -48,8 +48,7 @@
|
||||
"@backstage/config": "workspace:^",
|
||||
"@backstage/plugin-events-node": "workspace:^",
|
||||
"@backstage/types": "workspace:^",
|
||||
"luxon": "^3.0.0",
|
||||
"winston": "^3.2.1"
|
||||
"luxon": "^3.0.0"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@aws-sdk/types": "^3.347.0",
|
||||
|
||||
+16
-14
@@ -22,7 +22,7 @@ import {
|
||||
import { getVoidLogger } from '@backstage/backend-common';
|
||||
import { PluginTaskScheduler } from '@backstage/backend-tasks';
|
||||
import { ConfigReader } from '@backstage/config';
|
||||
import { TestEventBroker } from '@backstage/plugin-events-backend-test-utils';
|
||||
import { TestEventsService } from '@backstage/plugin-events-backend-test-utils';
|
||||
import { mockClient } from 'aws-sdk-client-mock';
|
||||
import { AwsSqsConsumingEventPublisher } from './AwsSqsConsumingEventPublisher';
|
||||
|
||||
@@ -53,12 +53,14 @@ describe('AwsSqsConsumingEventPublisher', () => {
|
||||
},
|
||||
});
|
||||
const logger = getVoidLogger();
|
||||
const events = new TestEventsService();
|
||||
const scheduler = {
|
||||
scheduleTask: jest.fn(),
|
||||
} as unknown as PluginTaskScheduler;
|
||||
|
||||
const publishers = AwsSqsConsumingEventPublisher.fromConfig({
|
||||
config,
|
||||
events,
|
||||
logger,
|
||||
scheduler,
|
||||
});
|
||||
@@ -85,21 +87,21 @@ describe('AwsSqsConsumingEventPublisher', () => {
|
||||
},
|
||||
});
|
||||
const logger = getVoidLogger();
|
||||
const events = new TestEventsService();
|
||||
const scheduler = {
|
||||
scheduleTask: jest.fn(),
|
||||
} as unknown as PluginTaskScheduler;
|
||||
|
||||
const publishers = AwsSqsConsumingEventPublisher.fromConfig({
|
||||
config,
|
||||
events,
|
||||
logger,
|
||||
scheduler,
|
||||
});
|
||||
expect(publishers.length).toEqual(1);
|
||||
|
||||
const publisher = publishers[0];
|
||||
|
||||
const eventBroker = new TestEventBroker();
|
||||
await publisher.setEventBroker(eventBroker);
|
||||
await publisher.start();
|
||||
|
||||
// publisher.connect(..) was causing the polling for events to be scheduled
|
||||
expect(scheduler.scheduleTask).toHaveBeenCalledWith(
|
||||
@@ -133,6 +135,7 @@ describe('AwsSqsConsumingEventPublisher', () => {
|
||||
},
|
||||
});
|
||||
const logger = getVoidLogger();
|
||||
const events = new TestEventsService();
|
||||
let taskFn: (() => Promise<void>) | undefined = undefined;
|
||||
const scheduler = {
|
||||
scheduleTask: (spec: { fn: () => Promise<void> }) => {
|
||||
@@ -196,32 +199,31 @@ describe('AwsSqsConsumingEventPublisher', () => {
|
||||
|
||||
const publishers = AwsSqsConsumingEventPublisher.fromConfig({
|
||||
config,
|
||||
events,
|
||||
logger,
|
||||
scheduler,
|
||||
});
|
||||
expect(publishers.length).toEqual(1);
|
||||
const publisher = publishers[0];
|
||||
|
||||
const eventBroker = new TestEventBroker();
|
||||
await publisher.setEventBroker(eventBroker);
|
||||
await publisher.start();
|
||||
|
||||
await taskFn!();
|
||||
await taskFn!();
|
||||
await taskFn!();
|
||||
|
||||
expect(eventBroker.published.length).toEqual(2);
|
||||
expect(eventBroker.published[0].topic).toEqual('fake1');
|
||||
expect(eventBroker.published[0].eventPayload).toEqual({
|
||||
expect(events.published).toHaveLength(2);
|
||||
expect(events.published[0].topic).toEqual('fake1');
|
||||
expect(events.published[0].eventPayload).toEqual({
|
||||
event: 'payload1',
|
||||
});
|
||||
expect(eventBroker.published[0].metadata).toEqual({
|
||||
expect(events.published[0].metadata).toEqual({
|
||||
'X-Custom-Attr': 'value',
|
||||
});
|
||||
|
||||
expect(eventBroker.published[1].topic).toEqual('fake1');
|
||||
expect(eventBroker.published[1].eventPayload).toEqual({
|
||||
expect(events.published[1].topic).toEqual('fake1');
|
||||
expect(events.published[1].eventPayload).toEqual({
|
||||
event: 'payload2',
|
||||
});
|
||||
expect(eventBroker.published[1].metadata).toEqual({});
|
||||
expect(events.published[1].metadata).toEqual({});
|
||||
});
|
||||
});
|
||||
|
||||
+15
-14
@@ -21,10 +21,10 @@ import {
|
||||
ReceiveMessageCommandInput,
|
||||
SQSClient,
|
||||
} from '@aws-sdk/client-sqs';
|
||||
import { LoggerService } from '@backstage/backend-plugin-api';
|
||||
import { PluginTaskScheduler } from '@backstage/backend-tasks';
|
||||
import { Config } from '@backstage/config';
|
||||
import { EventBroker, EventPublisher } from '@backstage/plugin-events-node';
|
||||
import { Logger } from 'winston';
|
||||
import { EventsService } from '@backstage/plugin-events-node';
|
||||
import { AwsSqsEventSourceConfig, readConfig } from './config';
|
||||
|
||||
/**
|
||||
@@ -34,28 +34,34 @@ import { AwsSqsEventSourceConfig, readConfig } from './config';
|
||||
* @public
|
||||
*/
|
||||
// TODO(pjungermann): add prom metrics? (see plugins/catalog-backend/src/util/metrics.ts, etc.)
|
||||
export class AwsSqsConsumingEventPublisher implements EventPublisher {
|
||||
export class AwsSqsConsumingEventPublisher {
|
||||
private readonly topic: string;
|
||||
private readonly receiveParams: ReceiveMessageCommandInput;
|
||||
private readonly sqs: SQSClient;
|
||||
private readonly queueUrl: string;
|
||||
private readonly taskTimeoutSeconds: number;
|
||||
private readonly waitTimeAfterEmptyReceiveMs;
|
||||
private eventBroker?: EventBroker;
|
||||
|
||||
static fromConfig(env: {
|
||||
config: Config;
|
||||
logger: Logger;
|
||||
events: EventsService;
|
||||
logger: LoggerService;
|
||||
scheduler: PluginTaskScheduler;
|
||||
}): AwsSqsConsumingEventPublisher[] {
|
||||
return readConfig(env.config).map(
|
||||
config =>
|
||||
new AwsSqsConsumingEventPublisher(env.logger, env.scheduler, config),
|
||||
new AwsSqsConsumingEventPublisher(
|
||||
env.logger,
|
||||
env.events,
|
||||
env.scheduler,
|
||||
config,
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
private constructor(
|
||||
private readonly logger: Logger,
|
||||
private readonly logger: LoggerService,
|
||||
private readonly events: EventsService,
|
||||
private readonly scheduler: PluginTaskScheduler,
|
||||
config: AwsSqsEventSourceConfig,
|
||||
) {
|
||||
@@ -80,12 +86,7 @@ export class AwsSqsConsumingEventPublisher implements EventPublisher {
|
||||
config.waitTimeAfterEmptyReceive.as('milliseconds');
|
||||
}
|
||||
|
||||
async setEventBroker(eventBroker: EventBroker): Promise<void> {
|
||||
this.eventBroker = eventBroker;
|
||||
return this.start();
|
||||
}
|
||||
|
||||
private async start(): Promise<void> {
|
||||
async start(): Promise<void> {
|
||||
const id = `events.awsSqs.publisher:${this.topic}`;
|
||||
const logger = this.logger.child({
|
||||
class: AwsSqsConsumingEventPublisher.prototype.constructor.name,
|
||||
@@ -172,7 +173,7 @@ export class AwsSqsConsumingEventPublisher implements EventPublisher {
|
||||
}
|
||||
});
|
||||
|
||||
this.eventBroker!.publish({
|
||||
this.events.publish({
|
||||
topic: this.topic,
|
||||
eventPayload,
|
||||
metadata,
|
||||
|
||||
+11
-17
@@ -14,26 +14,28 @@
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
import { createServiceFactory } from '@backstage/backend-plugin-api';
|
||||
import { mockServices, startTestBackend } from '@backstage/backend-test-utils';
|
||||
import { eventsExtensionPoint } from '@backstage/plugin-events-node/alpha';
|
||||
import { TestEventBroker } from '@backstage/plugin-events-backend-test-utils';
|
||||
import { eventsServiceRef } from '@backstage/plugin-events-node';
|
||||
import { TestEventsService } from '@backstage/plugin-events-backend-test-utils';
|
||||
import { eventsModuleAwsSqsConsumingEventPublisher } from './eventsModuleAwsSqsConsumingEventPublisher';
|
||||
import { AwsSqsConsumingEventPublisher } from '../publisher/AwsSqsConsumingEventPublisher';
|
||||
|
||||
describe('eventsModuleAwsSqsConsumingEventPublisher', () => {
|
||||
it('should be correctly wired and set up', async () => {
|
||||
let addedPublishers: AwsSqsConsumingEventPublisher[] | undefined;
|
||||
const extensionPoint = {
|
||||
addPublishers: (publishers: any) => {
|
||||
addedPublishers = publishers;
|
||||
const events = new TestEventsService();
|
||||
const eventsServiceFactory = createServiceFactory({
|
||||
service: eventsServiceRef,
|
||||
deps: {},
|
||||
async factory({}) {
|
||||
return events;
|
||||
},
|
||||
};
|
||||
});
|
||||
|
||||
const scheduler = mockServices.scheduler.mock();
|
||||
|
||||
await startTestBackend({
|
||||
extensionPoints: [[eventsExtensionPoint, extensionPoint]],
|
||||
features: [
|
||||
eventsServiceFactory(),
|
||||
eventsModuleAwsSqsConsumingEventPublisher(),
|
||||
mockServices.rootConfig.factory({
|
||||
data: {
|
||||
@@ -65,14 +67,6 @@ describe('eventsModuleAwsSqsConsumingEventPublisher', () => {
|
||||
],
|
||||
});
|
||||
|
||||
expect(addedPublishers).not.toBeUndefined();
|
||||
expect(addedPublishers!.length).toEqual(2);
|
||||
|
||||
const eventBroker = new TestEventBroker();
|
||||
await Promise.all(
|
||||
addedPublishers!.map(publisher => publisher.setEventBroker(eventBroker)),
|
||||
);
|
||||
|
||||
// publisher.connect(..) was causing the polling for events to be scheduled
|
||||
expect(scheduler.scheduleTask).toHaveBeenCalledWith(
|
||||
expect.objectContaining({ id: 'events.awsSqs.publisher:fake1' }),
|
||||
|
||||
+7
-8
@@ -18,8 +18,7 @@ import {
|
||||
coreServices,
|
||||
createBackendModule,
|
||||
} from '@backstage/backend-plugin-api';
|
||||
import { loggerToWinstonLogger } from '@backstage/backend-common';
|
||||
import { eventsExtensionPoint } from '@backstage/plugin-events-node/alpha';
|
||||
import { eventsServiceRef } from '@backstage/plugin-events-node';
|
||||
import { AwsSqsConsumingEventPublisher } from '../publisher/AwsSqsConsumingEventPublisher';
|
||||
|
||||
/**
|
||||
@@ -34,19 +33,19 @@ export const eventsModuleAwsSqsConsumingEventPublisher = createBackendModule({
|
||||
env.registerInit({
|
||||
deps: {
|
||||
config: coreServices.rootConfig,
|
||||
events: eventsExtensionPoint,
|
||||
events: eventsServiceRef,
|
||||
logger: coreServices.logger,
|
||||
scheduler: coreServices.scheduler,
|
||||
},
|
||||
async init({ config, events, logger, scheduler }) {
|
||||
const winstonLogger = loggerToWinstonLogger(logger);
|
||||
const sqs = AwsSqsConsumingEventPublisher.fromConfig({
|
||||
config: config,
|
||||
logger: winstonLogger,
|
||||
scheduler: scheduler,
|
||||
config,
|
||||
events,
|
||||
logger,
|
||||
scheduler,
|
||||
});
|
||||
|
||||
events.addPublishers(sqs);
|
||||
await Promise.all(sqs.map(publisher => publisher.start()));
|
||||
},
|
||||
});
|
||||
},
|
||||
|
||||
Reference in New Issue
Block a user