MCPcopy Create free account
hub / github.com/cameri/nostream / SubscribeMessageHandler

Class SubscribeMessageHandler

src/handlers/subscribe-message-handler.ts:24–134  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

22const logger = createLogger('subscribe-message-handler')
23
24export 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),

Callers

nothing calls this directly

Calls

no outgoing calls

Tested by

no test coverage detected