Skip to content

Commit 2823712

Browse files
authored
Merge branch 'main' into ADMIN-CONSOLE-2-SSE
2 parents f052b55 + 237b1a4 commit 2823712

13 files changed

Lines changed: 68 additions & 35 deletions
Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
"nostream": minor
3+
---
4+
5+
feat: instrument event and websocket handlers with OpenTelemetry metrics

src/adapters/web-socket-adapter.ts

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ import { WebSocketAdapterEvent, WebSocketServerAdapterEvent } from '../constants
1515
import { attemptValidation } from '../utils/validation'
1616
import { ContextMetadataKey } from '../constants/base'
1717
import { createLogger } from '../factories/logger-factory'
18+
import { recordWebsocketConnectionClosed, recordWebsocketConnectionOpened } from '../telemetry/event-metrics'
1819
import { Event } from '../@types/event'
1920
import { getRemoteAddress } from '../utils/http'
2021
import { IRateLimiter } from '../@types/utils'
@@ -82,6 +83,7 @@ export class WebSocketAdapter extends EventEmitter implements IWebSocketAdapter
8283
.on(WebSocketAdapterEvent.Message, this.sendMessage.bind(this))
8384

8485
logger('client %s connected from %s', this.clientId, this.clientAddress.address)
86+
recordWebsocketConnectionOpened()
8587

8688
// NIP-42
8789
this.challenge = randomBytes(32).toString('base64url')
@@ -260,6 +262,7 @@ export class WebSocketAdapter extends EventEmitter implements IWebSocketAdapter
260262
}
261263

262264
private onClientClose() {
265+
recordWebsocketConnectionClosed()
263266
this.alive = false
264267
this.subscriptions.clear()
265268
this.authenticatedPubkeys.clear()

src/handlers/event-message-handler.ts

Lines changed: 14 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,7 @@ import {
3131
import { IEventRepository, INip05VerificationRepository, IUserRepository } from '../@types/repositories'
3232
import { IEventStrategy, IMessageHandler } from '../@types/message-handlers'
3333
import { CacheAdmissionState } from '../constants/caching'
34-
import { createCommandResult } from '../utils/messages'
34+
import { createEventCommandResult } from '../telemetry/event-metrics'
3535
import { createLogger } from '../factories/logger-factory'
3636
import { Factory } from '../@types/base'
3737
import { ICacheAdapter } from '../@types/adapters'
@@ -63,13 +63,13 @@ export class EventMessageHandler implements IMessageHandler {
6363
let reason = await this.isEventValid(event)
6464
if (reason) {
6565
logger('event %s rejected: %s', event.id, reason)
66-
this.webSocket.emit(WebSocketAdapterEvent.Message, createCommandResult(event.id, false, reason))
66+
this.webSocket.emit(WebSocketAdapterEvent.Message, createEventCommandResult(event.id, false, reason))
6767
return
6868
}
6969

7070
if (isExpiredEvent(event)) {
7171
logger('event %s rejected: expired')
72-
this.webSocket.emit(WebSocketAdapterEvent.Message, createCommandResult(event.id, false, 'event is expired'))
72+
this.webSocket.emit(WebSocketAdapterEvent.Message, createEventCommandResult(event.id, false, 'event is expired'))
7373
return
7474
}
7575

@@ -79,43 +79,43 @@ export class EventMessageHandler implements IMessageHandler {
7979
logger('event %s rejected: rate-limited')
8080
this.webSocket.emit(
8181
WebSocketAdapterEvent.Message,
82-
createCommandResult(event.id, false, 'rate-limited: slow down'),
82+
createEventCommandResult(event.id, false, 'rate-limited: slow down'),
8383
)
8484
return
8585
}
8686

8787
reason = this.canAcceptEvent(event)
8888
if (reason) {
8989
logger('event %s rejected: %s', event.id, reason)
90-
this.webSocket.emit(WebSocketAdapterEvent.Message, createCommandResult(event.id, false, reason))
90+
this.webSocket.emit(WebSocketAdapterEvent.Message, createEventCommandResult(event.id, false, reason))
9191
return
9292
}
9393

9494
reason = await this.isProtectedEventBlocked(event)
9595
if (reason) {
9696
logger('event %s rejected: %s', event.id, reason)
97-
this.webSocket.emit(WebSocketAdapterEvent.Message, createCommandResult(event.id, false, reason))
97+
this.webSocket.emit(WebSocketAdapterEvent.Message, createEventCommandResult(event.id, false, reason))
9898
return
9999
}
100100

101101
reason = await this.isBlockedByRequestToVanish(event)
102102
if (reason) {
103103
logger('event %s rejected: %s', event.id, reason)
104-
this.webSocket.emit(WebSocketAdapterEvent.Message, createCommandResult(event.id, false, reason))
104+
this.webSocket.emit(WebSocketAdapterEvent.Message, createEventCommandResult(event.id, false, reason))
105105
return
106106
}
107107

108108
reason = await this.isUserAdmitted(event)
109109
if (reason) {
110110
logger('event %s rejected: %s', event.id, reason)
111-
this.webSocket.emit(WebSocketAdapterEvent.Message, createCommandResult(event.id, false, reason))
111+
this.webSocket.emit(WebSocketAdapterEvent.Message, createEventCommandResult(event.id, false, reason))
112112
return
113113
}
114114

115115
reason = await this.checkNip05Verification(event)
116116
if (reason) {
117117
logger('event %s rejected: %s', event.id, reason)
118-
this.webSocket.emit(WebSocketAdapterEvent.Message, createCommandResult(event.id, false, reason))
118+
this.webSocket.emit(WebSocketAdapterEvent.Message, createEventCommandResult(event.id, false, reason))
119119
return
120120
}
121121

@@ -124,7 +124,7 @@ export class EventMessageHandler implements IMessageHandler {
124124
if (typeof strategy?.execute !== 'function') {
125125
this.webSocket.emit(
126126
WebSocketAdapterEvent.Message,
127-
createCommandResult(event.id, false, 'error: event not supported'),
127+
createEventCommandResult(event.id, false, 'error: event not supported'),
128128
)
129129
return
130130
}
@@ -134,7 +134,7 @@ export class EventMessageHandler implements IMessageHandler {
134134
this.processNip05Metadata(event)
135135
} catch (error) {
136136
logger.error('error handling message', message, error)
137-
this.webSocket.emit(WebSocketAdapterEvent.Message, createCommandResult(event.id, false, 'error: unable to process event'))
137+
this.webSocket.emit(WebSocketAdapterEvent.Message, createEventCommandResult(event.id, false, 'error: unable to process event'))
138138
}
139139
}
140140

@@ -242,7 +242,9 @@ export class EventMessageHandler implements IMessageHandler {
242242
}
243243

244244
const checkEmbedded = async (evt: Event, depth = 0): Promise<boolean> => {
245-
if (depth > 10) return false // Prevent infinite loops or excessive recursion
245+
if (depth > 10) {
246+
return false // Prevent infinite loops or excessive recursion
247+
}
246248
if ((evt.kind === EventKinds.REPOST || evt.kind === EventKinds.GENERIC_REPOST) && evt.content.length > 0) {
247249
try {
248250
const embedded = attemptValidation(eventSchema)(JSON.parse(evt.content)) as Event

src/handlers/event-strategies/default-event-strategy.ts

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
import { createCommandResult } from '../../utils/messages'
1+
import { createEventCommandResult } from '../../telemetry/event-metrics'
22
import { createLogger } from '../../factories/logger-factory'
33
import { Event } from '../../@types/event'
44
import { IEventRepository } from '../../@types/repositories'
@@ -17,7 +17,7 @@ export class DefaultEventStrategy implements IEventStrategy<Event, Promise<void>
1717
public async execute(event: Event): Promise<void> {
1818
logger('received event: %o', event)
1919
const count = await this.eventRepository.create(event)
20-
this.webSocket.emit(WebSocketAdapterEvent.Message, createCommandResult(event.id, true, count ? '' : 'duplicate:'))
20+
this.webSocket.emit(WebSocketAdapterEvent.Message, createEventCommandResult(event.id, true, count ? '' : 'duplicate:'))
2121

2222
if (count) {
2323
this.webSocket.emit(WebSocketAdapterEvent.Broadcast, event)

src/handlers/event-strategies/delete-event-strategy.ts

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
import { createCommandResult } from '../../utils/messages'
1+
import { createEventCommandResult } from '../../telemetry/event-metrics'
22
import { createLogger } from '../../factories/logger-factory'
33
import { Event } from '../../@types/event'
44
import { EventTags } from '../../constants/base'
@@ -31,7 +31,7 @@ export class DeleteEventStrategy implements IEventStrategy<Event, Promise<void>>
3131
}
3232

3333
const count = await this.eventRepository.create(event)
34-
this.webSocket.emit(WebSocketAdapterEvent.Message, createCommandResult(event.id, true, count ? '' : 'duplicate:'))
34+
this.webSocket.emit(WebSocketAdapterEvent.Message, createEventCommandResult(event.id, true, count ? '' : 'duplicate:'))
3535

3636
if (count) {
3737
this.webSocket.emit(WebSocketAdapterEvent.Broadcast, event)

src/handlers/event-strategies/ephemeral-event-strategy.ts

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
import { createCommandResult } from '../../utils/messages'
1+
import { createEventCommandResult } from '../../telemetry/event-metrics'
22
import { createLogger } from '../../factories/logger-factory'
33
import { Event } from '../../@types/event'
44
import { IEventStrategy } from '../../@types/message-handlers'
@@ -12,7 +12,7 @@ export class EphemeralEventStrategy implements IEventStrategy<Event, Promise<voi
1212

1313
public async execute(event: Event): Promise<void> {
1414
logger('received ephemeral event: %o', event)
15-
this.webSocket.emit(WebSocketAdapterEvent.Message, createCommandResult(event.id, true, ''))
15+
this.webSocket.emit(WebSocketAdapterEvent.Message, createEventCommandResult(event.id, true, ''))
1616
this.webSocket.emit(WebSocketAdapterEvent.Broadcast, event)
1717
}
1818
}

src/handlers/event-strategies/gift-wrap-event-strategy.ts

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
import { createCommandResult } from '../../utils/messages'
1+
import { createEventCommandResult } from '../../telemetry/event-metrics'
22
import { createLogger } from '../../factories/logger-factory'
33
import { Event } from '../../@types/event'
44
import { EventTags } from '../../constants/base'
@@ -21,12 +21,12 @@ export class GiftWrapEventStrategy implements IEventStrategy<Event, Promise<void
2121

2222
const reason = this.validateGiftWrap(event)
2323
if (reason) {
24-
this.webSocket.emit(WebSocketAdapterEvent.Message, createCommandResult(event.id, false, `invalid: ${reason}`))
24+
this.webSocket.emit(WebSocketAdapterEvent.Message, createEventCommandResult(event.id, false, `invalid: ${reason}`))
2525
return
2626
}
2727

2828
const count = await this.eventRepository.create(event)
29-
this.webSocket.emit(WebSocketAdapterEvent.Message, createCommandResult(event.id, true, count ? '' : 'duplicate:'))
29+
this.webSocket.emit(WebSocketAdapterEvent.Message, createEventCommandResult(event.id, true, count ? '' : 'duplicate:'))
3030

3131
if (count) {
3232
this.webSocket.emit(WebSocketAdapterEvent.Broadcast, event)

src/handlers/event-strategies/group-event-strategy.ts

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
import { createCommandResult } from '../../utils/messages'
1+
import { createEventCommandResult } from '../../telemetry/event-metrics'
22
import { createLogger } from '../../factories/logger-factory'
33
import { Event } from '../../@types/event'
44
import { EventTags } from '../../constants/base'
@@ -20,12 +20,12 @@ export class GroupEventStrategy implements IEventStrategy<Event, Promise<void>>
2020

2121
const reason = this.validateGroupEvent(event)
2222
if (reason) {
23-
this.webSocket.emit(WebSocketAdapterEvent.Message, createCommandResult(event.id, false, `invalid: ${reason}`))
23+
this.webSocket.emit(WebSocketAdapterEvent.Message, createEventCommandResult(event.id, false, `invalid: ${reason}`))
2424
return
2525
}
2626

2727
const count = await this.eventRepository.create(event)
28-
this.webSocket.emit(WebSocketAdapterEvent.Message, createCommandResult(event.id, true, count ? '' : 'duplicate:'))
28+
this.webSocket.emit(WebSocketAdapterEvent.Message, createEventCommandResult(event.id, true, count ? '' : 'duplicate:'))
2929

3030
if (count) {
3131
this.webSocket.emit(WebSocketAdapterEvent.Broadcast, event)

src/handlers/event-strategies/parameterized-replaceable-event-strategy.ts

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
import { Event, ParameterizedReplaceableEvent } from '../../@types/event'
22
import { EventDeduplicationMetadataKey, EventTags } from '../../constants/base'
3-
import { createCommandResult } from '../../utils/messages'
3+
import { createEventCommandResult } from '../../telemetry/event-metrics'
44
import { createLogger } from '../../factories/logger-factory'
55
import { IEventRepository } from '../../@types/repositories'
66
import { IEventStrategy } from '../../@types/message-handlers'
@@ -29,7 +29,7 @@ export class ParameterizedReplaceableEventStrategy implements IEventStrategy<Eve
2929
}
3030

3131
const count = await this.eventRepository.upsert(parameterizedReplaceableEvent)
32-
this.webSocket.emit(WebSocketAdapterEvent.Message, createCommandResult(event.id, true, count ? '' : 'duplicate:'))
32+
this.webSocket.emit(WebSocketAdapterEvent.Message, createEventCommandResult(event.id, true, count ? '' : 'duplicate:'))
3333

3434
if (count) {
3535
this.webSocket.emit(WebSocketAdapterEvent.Broadcast, event)

src/handlers/event-strategies/replaceable-event-strategy.ts

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
import { createCommandResult } from '../../utils/messages'
1+
import { createEventCommandResult } from '../../telemetry/event-metrics'
22
import { createLogger } from '../../factories/logger-factory'
33
import { Event } from '../../@types/event'
44
import { IEventRepository } from '../../@types/repositories'
@@ -18,7 +18,7 @@ export class ReplaceableEventStrategy implements IEventStrategy<Event, Promise<v
1818
logger('received replaceable event: %o', event)
1919
try {
2020
const count = await this.eventRepository.upsert(event)
21-
this.webSocket.emit(WebSocketAdapterEvent.Message, createCommandResult(event.id, true, count ? '' : 'duplicate:'))
21+
this.webSocket.emit(WebSocketAdapterEvent.Message, createEventCommandResult(event.id, true, count ? '' : 'duplicate:'))
2222
if (count) {
2323
this.webSocket.emit(WebSocketAdapterEvent.Broadcast, event)
2424
}
@@ -27,12 +27,12 @@ export class ReplaceableEventStrategy implements IEventStrategy<Event, Promise<v
2727
if (error.message.endsWith('duplicate key value violates unique constraint "events_event_id_unique"')) {
2828
this.webSocket.emit(
2929
WebSocketAdapterEvent.Message,
30-
createCommandResult(event.id, false, 'rejected: event already exists'),
30+
createEventCommandResult(event.id, false, 'rejected: event already exists'),
3131
)
3232
return
3333
}
3434

35-
this.webSocket.emit(WebSocketAdapterEvent.Message, createCommandResult(event.id, false, 'error: '))
35+
this.webSocket.emit(WebSocketAdapterEvent.Message, createEventCommandResult(event.id, false, 'error: '))
3636
}
3737
}
3838
}

0 commit comments

Comments
 (0)