import type { BaseEvent } from './Events' import type { AgentContext } from './context' import type { EventEmitter as NativeEventEmitter } from 'events' import { fromEventPattern, Subject } from 'rxjs' import { takeUntil } from 'rxjs/operators' import { InjectionSymbols } from '../constants' import { injectable, inject } from '../plugins' import { AgentDependencies } from './AgentDependencies' type EmitEvent = Omit @injectable() export class EventEmitter { private eventEmitter: NativeEventEmitter private stop$: Subject public constructor( @inject(InjectionSymbols.AgentDependencies) agentDependencies: AgentDependencies, @inject(InjectionSymbols.Stop$) stop$: Subject ) { this.eventEmitter = new agentDependencies.EventEmitterClass() this.stop$ = stop$ } // agentContext is currently not used, but already making required as it will be used soon public emit(agentContext: AgentContext, data: EmitEvent) { this.eventEmitter.emit(data.type, { ...data, metadata: { contextCorrelationId: agentContext.contextCorrelationId, }, }) } public on(event: T['type'], listener: (data: T) => void | Promise) { this.eventEmitter.on(event, listener) } public off(event: T['type'], listener: (data: T) => void | Promise) { this.eventEmitter.off(event, listener) } public observable(event: T['type']) { return fromEventPattern( (handler) => this.on(event, handler), (handler) => this.off(event, handler) ).pipe(takeUntil(this.stop$)) } }