This commit is contained in:
JupiterWalker 2025-03-06 05:40:43 +01:00 committed by GitHub
commit 530d6666ec
No known key found for this signature in database
GPG key ID: B5690EEEBB952194

View file

@ -99,7 +99,7 @@ export class MessageEventBus extends EventEmitter {
for (const destinationData of savedEventDestinations) {
try {
const destination = messageEventBusDestinationFromDb(this, destinationData);
if (destination) {
if (destination && destination.enabled) {
await this.addDestination(destination, false);
}
} catch (error) {
@ -206,21 +206,36 @@ export class MessageEventBus extends EventEmitter {
}
async addDestination(destination: MessageEventBusDestination, notifyWorkers: boolean = true) {
await this.removeDestination(destination.getId(), false);
this.destinations[destination.getId()] = destination;
this.destinations[destination.getId()].startListening();
if (notifyWorkers) {
await this.removeDestination(destination.getId(), false);
await destination.saveToDb();
void this.publisher.publishCommand({ command: 'restart-event-bus' });
}
this.destinations[destination.getId()] = destination;
this.destinations[destination.getId()].startListening();
return destination;
}
async findDestination(id?: string): Promise<MessageEventBusDestinationOptions[]> {
let result: MessageEventBusDestinationOptions[];
if (id && Object.keys(this.destinations).includes(id)) {
result = [this.destinations[id].serialize()];
} else {
result = Object.keys(this.destinations).map((e) => this.destinations[e].serialize());
let result: MessageEventBusDestinationOptions[] = [];
if (id) {
const savedEventDestination = await this.eventDestinationsRepository.findOne({ where: { id } });
if (savedEventDestination) {
const destination = messageEventBusDestinationFromDb(this, savedEventDestination);
if (destination) {
result = [destination.serialize()];
}
}
} else {
const savedEventDestinations = await this.eventDestinationsRepository.find();
if (savedEventDestinations.length > 0) {
for (const destinationData of savedEventDestinations) {
const destination = messageEventBusDestinationFromDb(this, destinationData)
if (destination) {
result.push(destination.serialize());
}
}
}
}
return result.sort((a, b) => (a.__type ?? '').localeCompare(b.__type ?? ''));
}