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:
Patrick Jungermann
2024-01-23 21:01:42 +01:00
parent 8f6afa94a9
commit 132d672747
8 changed files with 88 additions and 62 deletions
+32
View File
@@ -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",
@@ -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({});
});
});
@@ -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,
@@ -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' }),
@@ -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()));
},
});
},
-1
View File
@@ -6391,7 +6391,6 @@ __metadata:
"@backstage/types": "workspace:^"
aws-sdk-client-mock: ^3.0.0
luxon: ^3.0.0
winston: ^3.2.1
languageName: unknown
linkType: soft