| 22 | const logger = createLogger('subscribe-message-handler') |
| 23 | |
| 24 | export class SubscribeMessageHandler implements IMessageHandler, IAbortable { |
| 25 | //private readonly abortController: AbortController |
| 26 | |
| 27 | public constructor( |
| 28 | private readonly webSocket: IWebSocketAdapter, |
| 29 | private readonly eventRepository: IEventRepository, |
| 30 | private readonly settings: () => Settings, |
| 31 | ) { |
| 32 | //this.abortController = new AbortController() |
| 33 | } |
| 34 | |
| 35 | public abort(): void { |
| 36 | //this.abortController.abort() |
| 37 | } |
| 38 | |
| 39 | public async handleMessage(message: SubscribeMessage): Promise<void> { |
| 40 | const subscriptionId = message[1] |
| 41 | const filters = uniqWith(equals, message.slice(2)) as SubscriptionFilter[] |
| 42 | |
| 43 | const reason = this.canSubscribe(subscriptionId, filters) |
| 44 | if (reason) { |
| 45 | logger('subscription %s with %o rejected: %s', subscriptionId, filters, reason) |
| 46 | this.webSocket.emit(WebSocketAdapterEvent.Message, createNoticeMessage(`Subscription rejected: ${reason}`)) |
| 47 | return |
| 48 | } |
| 49 | |
| 50 | this.webSocket.emit(WebSocketAdapterEvent.Subscribe, subscriptionId, filters) |
| 51 | |
| 52 | await this.fetchAndSend(subscriptionId, filters) |
| 53 | } |
| 54 | |
| 55 | private async fetchAndSend(subscriptionId: string, filters: SubscriptionFilter[]): Promise<void> { |
| 56 | logger('fetching events for subscription %s with filters %o', subscriptionId, filters) |
| 57 | const sendEvent = (event: Event) => |
| 58 | this.webSocket.emit(WebSocketAdapterEvent.Message, createOutgoingEventMessage(subscriptionId, event)) |
| 59 | const sendEOSE = () => |
| 60 | this.webSocket.emit(WebSocketAdapterEvent.Message, createEndOfStoredEventsNoticeMessage(subscriptionId)) |
| 61 | const isSubscribedToEvent = SubscribeMessageHandler.isClientSubscribedToEvent(filters) |
| 62 | const isTagUnexpired = (event: Event) => { |
| 63 | if (isExpiredEvent(event)) { |
| 64 | return false |
| 65 | } |
| 66 | return true |
| 67 | } |
| 68 | |
| 69 | const findEvents = this.eventRepository.findByFilters(filters).stream() |
| 70 | |
| 71 | // const abortableFindEvents = addAbortSignal(this.abortController.signal, findEvents) |
| 72 | |
| 73 | try { |
| 74 | await pipeline( |
| 75 | findEvents, |
| 76 | streamFilter(propSatisfies(isNil, 'deleted_at')), |
| 77 | streamMap(toNostrEvent), |
| 78 | streamFilter(isTagUnexpired), |
| 79 | streamFilter(isSubscribedToEvent), |
| 80 | streamEach(sendEvent), |
| 81 | streamEnd(sendEOSE), |
nothing calls this directly
no outgoing calls
no test coverage detected