Skip to content

Commit 275c30c

Browse files
feat(dvm): trap NIP-90 job request events and record them via job repository (#729)
* feat(dvm): trap NIP-90 job request events and record them via job repository Signed-off-by: Priyanshubhartistm <bhartipriyanshustm@gmail.com> * test(dvm): add tests for dvm job request ingestion Signed-off-by: Priyanshubhartistm <bhartipriyanshustm@gmail.com> --------- Signed-off-by: Priyanshubhartistm <bhartipriyanshustm@gmail.com>
1 parent 9fb1c10 commit 275c30c

14 files changed

Lines changed: 331 additions & 29 deletions

‎.changeset/dvm-job-ingestion.md‎

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(dvm): trap NIP-90 job request events (kind 5000-5999) and record them via the job repository

‎.knip.json‎

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,6 @@
1616
],
1717
"ignore": [
1818
".nostr/**",
19-
"src/repositories/dvm-job-repository.ts",
2019
"src/utils/relay-probe/**"
2120
],
2221
"commitlint": false,

‎src/constants/base.ts‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,9 @@ export enum EventKinds {
4141
// Lightning zaps
4242
ZAP_REQUEST = 9734,
4343
ZAP_RECEIPT = 9735,
44+
// NIP-90: Data Vending Machines — job request events
45+
DVM_JOB_REQUEST_FIRST = 5000,
46+
DVM_JOB_REQUEST_LAST = 5999,
4447
// Replaceable events
4548
REPLACEABLE_FIRST = 10000,
4649
// NIP-65: Relay List Metadata

‎src/factories/event-strategy-factory.ts‎

Lines changed: 12 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,8 @@
11
import { ICacheAdapter, IWebSocketAdapter } from '../@types/adapters'
2-
import { IEventRepository, IInviteCodeRepository, IUserRepository } from '../@types/repositories'
2+
import { IDvmJobRepository, IEventRepository, IInviteCodeRepository, IUserRepository } from '../@types/repositories'
33
import {
44
isDeleteEvent,
5+
isDvmJobRequestEvent,
56
isEphemeralEvent,
67
isGiftWrapEvent,
78
isMarmotGroupEvent,
@@ -14,6 +15,7 @@ import { isNip43JoinRequest, isNip43LeaveRequest } from '../utils/nip43'
1415
import { isRelayListEvent } from '../utils/nip65'
1516
import { DefaultEventStrategy } from '../handlers/event-strategies/default-event-strategy'
1617
import { DeleteEventStrategy } from '../handlers/event-strategies/delete-event-strategy'
18+
import { DvmJobRequestEventStrategy } from '../handlers/event-strategies/dvm-job-request-event-strategy'
1719
import { EphemeralEventStrategy } from '../handlers/event-strategies/ephemeral-event-strategy'
1820
import { Event } from '../@types/event'
1921
import { Factory } from '../@types/base'
@@ -33,6 +35,7 @@ export const eventStrategyFactory =
3335
eventRepository: IEventRepository,
3436
userRepository: IUserRepository,
3537
inviteCodeRepository: IInviteCodeRepository,
38+
dvmJobRepository: IDvmJobRepository,
3639
cache: ICacheAdapter,
3740
settings: () => Settings,
3841
): Factory<IEventStrategy<Event, Promise<void>>, [Event, IWebSocketAdapter]> =>
@@ -47,12 +50,17 @@ export const eventStrategyFactory =
4750
return new TimestampEventStrategy(adapter, eventRepository)
4851
} else if (isRelayListEvent(event) || isReplaceableEvent(event)) {
4952
return new ReplaceableEventStrategy(adapter, eventRepository)
50-
// NIP-43: Join/Leave requests MUST be checked before the generic ephemeral
51-
// handler, because kinds 28934/28936 fall in the ephemeral range (20000-29999).
53+
// NIP-43: Join/Leave requests MUST be checked before the generic ephemeral
54+
// handler, because kinds 28934/28936 fall in the ephemeral range (20000-29999).
5255
} else if (isNip43JoinRequest(event)) {
5356
return new JoinRequestEventStrategy(adapter, inviteCodeRepository, userRepository, cache, settings)
5457
} else if (isNip43LeaveRequest(event)) {
5558
return new LeaveRequestEventStrategy(adapter, userRepository, cache, settings)
59+
// NIP-90: DVM job requests (kind 5000-5999) checked early, same reasoning
60+
// as the NIP-43 checks above — kept explicit rather than relying on it
61+
// falling through to DefaultEventStrategy.
62+
} else if (isDvmJobRequestEvent(event)) {
63+
return new DvmJobRequestEventStrategy(adapter, eventRepository, dvmJobRepository)
5664
} else if (isEphemeralEvent(event)) {
5765
return new EphemeralEventStrategy(adapter)
5866
} else if (isDeleteEvent(event)) {
@@ -62,4 +70,4 @@ export const eventStrategyFactory =
6270
}
6371

6472
return new DefaultEventStrategy(adapter, eventRepository)
65-
}
73+
}

‎src/factories/message-handler-factory.ts‎

Lines changed: 16 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,11 @@
11
import { ICacheAdapter, IWebSocketAdapter } from '../@types/adapters'
2-
import { IEventRepository, IInviteCodeRepository, INip05VerificationRepository, IUserRepository } from '../@types/repositories'
2+
import {
3+
IDvmJobRepository,
4+
IEventRepository,
5+
IInviteCodeRepository,
6+
INip05VerificationRepository,
7+
IUserRepository,
8+
} from '../@types/repositories'
39
import { IncomingMessage, MessageType } from '../@types/messages'
410
import { createSettings } from './settings-factory'
511
import { AuthMessageHandler } from '../handlers/auth-message-handler'
@@ -26,13 +32,21 @@ export const messageHandlerFactory =
2632
userRepository: IUserRepository,
2733
nip05VerificationRepository: INip05VerificationRepository,
2834
inviteCodeRepository: IInviteCodeRepository,
35+
dvmJobRepository: IDvmJobRepository,
2936
) =>
3037
([message, adapter]: [IncomingMessage, IWebSocketAdapter]) => {
3138
switch (message[0]) {
3239
case MessageType.EVENT: {
3340
return new EventMessageHandler(
3441
adapter,
35-
eventStrategyFactory(eventRepository, userRepository, inviteCodeRepository, getCache(), createSettings),
42+
eventStrategyFactory(
43+
eventRepository,
44+
userRepository,
45+
inviteCodeRepository,
46+
dvmJobRepository,
47+
getCache(),
48+
createSettings,
49+
),
3650
eventRepository,
3751
userRepository,
3852
createSettings,

‎src/factories/websocket-adapter-factory.ts‎

Lines changed: 15 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,13 @@
11
import { IncomingMessage } from 'http'
22
import { WebSocket } from 'ws'
33

4-
import { IEventRepository, IInviteCodeRepository, INip05VerificationRepository, IUserRepository } from '../@types/repositories'
4+
import {
5+
IDvmJobRepository,
6+
IEventRepository,
7+
IInviteCodeRepository,
8+
INip05VerificationRepository,
9+
IUserRepository,
10+
} from '../@types/repositories'
511
import { createSettings } from './settings-factory'
612
import { IWebSocketServerAdapter } from '../@types/adapters'
713
import { messageHandlerFactory } from './message-handler-factory'
@@ -14,13 +20,20 @@ export const webSocketAdapterFactory =
1420
userRepository: IUserRepository,
1521
nip05VerificationRepository: INip05VerificationRepository,
1622
inviteCodeRepository: IInviteCodeRepository,
23+
dvmJobRepository: IDvmJobRepository,
1724
) =>
1825
([client, request, webSocketServerAdapter]: [WebSocket, IncomingMessage, IWebSocketServerAdapter]) =>
1926
new WebSocketAdapter(
2027
client,
2128
request,
2229
webSocketServerAdapter,
23-
messageHandlerFactory(eventRepository, userRepository, nip05VerificationRepository, inviteCodeRepository),
30+
messageHandlerFactory(
31+
eventRepository,
32+
userRepository,
33+
nip05VerificationRepository,
34+
inviteCodeRepository,
35+
dvmJobRepository,
36+
),
2437
rateLimiterFactory,
2538
createSettings,
2639
)

‎src/factories/worker-factory.ts‎

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ import { AppWorker } from '../app/worker'
88
import { createLogger } from './logger-factory'
99
import { createSettings } from '../factories/settings-factory'
1010
import { createWebApp } from './web-app-factory'
11+
import { DvmJobRepository } from '../repositories/dvm-job-repository'
1112
import { EventRepository } from '../repositories/event-repository'
1213
import { InviteCodeRepository } from '../repositories/invite-code-repository'
1314
import { Nip05VerificationRepository } from '../repositories/nip05-verification-repository'
@@ -24,6 +25,7 @@ export const workerFactory = (): AppWorker => {
2425
const userRepository = new UserRepository(dbClient, eventRepository)
2526
const nip05VerificationRepository = new Nip05VerificationRepository(dbClient)
2627
const inviteCodeRepository = new InviteCodeRepository(dbClient)
28+
const dvmJobRepository = new DvmJobRepository(dbClient)
2729

2830
const settings = createSettings()
2931

@@ -65,7 +67,13 @@ export const workerFactory = (): AppWorker => {
6567
const adapter = new WebSocketServerAdapter(
6668
server,
6769
webSocketServer,
68-
webSocketAdapterFactory(eventRepository, userRepository, nip05VerificationRepository, inviteCodeRepository),
70+
webSocketAdapterFactory(
71+
eventRepository,
72+
userRepository,
73+
nip05VerificationRepository,
74+
inviteCodeRepository,
75+
dvmJobRepository,
76+
),
6977
createSettings,
7078
)
7179

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,42 @@
1+
import { createEventCommandResult } from '../../telemetry/event-metrics'
2+
import { createLogger } from '../../factories/logger-factory'
3+
import { Event } from '../../@types/event'
4+
import { IDvmJobRepository, IEventRepository } from '../../@types/repositories'
5+
import { IEventStrategy } from '../../@types/message-handlers'
6+
import { IWebSocketAdapter } from '../../@types/adapters'
7+
import { WebSocketAdapterEvent } from '../../constants/adapter'
8+
9+
const logger = createLogger('dvm-job-request-event-strategy')
10+
11+
export class DvmJobRequestEventStrategy implements IEventStrategy<Event, Promise<void>> {
12+
public constructor(
13+
private readonly webSocket: IWebSocketAdapter,
14+
private readonly eventRepository: IEventRepository,
15+
private readonly dvmJobRepository: IDvmJobRepository,
16+
) {}
17+
18+
public async execute(event: Event): Promise<void> {
19+
logger('received dvm job request: %o', event)
20+
21+
const count = await this.eventRepository.create(event)
22+
this.webSocket.emit(
23+
WebSocketAdapterEvent.Message,
24+
createEventCommandResult(event.id, true, count ? '' : 'duplicate:'),
25+
)
26+
27+
if (!count) {
28+
return
29+
}
30+
31+
this.webSocket.emit(WebSocketAdapterEvent.Broadcast, event)
32+
33+
try {
34+
await this.dvmJobRepository.create(event.id, event.pubkey, event.kind)
35+
} catch (error) {
36+
// Job-state recording is best-effort: the event itself is already
37+
// stored and broadcast correctly, so a repository failure here must
38+
// not surface as a rejection of a valid event.
39+
logger.error('unable to record dvm job for event %s: %o', event.id, error)
40+
}
41+
}
42+
}

‎src/utils/event.ts‎

Lines changed: 19 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -40,18 +40,20 @@ export const isEventMatchingFilter =
4040
(filter: SubscriptionFilter) =>
4141
(event: Event): boolean => {
4242
const startsWith = (input: string) => (prefix: string) => input.startsWith(prefix)
43-
const isMatchingGenericTagCriterion = (key: string, criterion: string) => (tag: Tag): boolean => {
44-
const [, tagName] = key
45-
if (tag[0] !== tagName) {
46-
return false
47-
}
43+
const isMatchingGenericTagCriterion =
44+
(key: string, criterion: string) =>
45+
(tag: Tag): boolean => {
46+
const [, tagName] = key
47+
if (tag[0] !== tagName) {
48+
return false
49+
}
4850

49-
if (isGeohashPrefixCriterion(key, criterion)) {
50-
return tag[1].startsWith(stripGeohashPrefixWildcard(criterion))
51-
}
51+
if (isGeohashPrefixCriterion(key, criterion)) {
52+
return tag[1].startsWith(stripGeohashPrefixWildcard(criterion))
53+
}
5254

53-
return tag[1] === criterion
54-
}
55+
return tag[1] === criterion
56+
}
5557

5658
// NIP-01: Basic protocol flow description
5759

@@ -96,7 +98,9 @@ export const isEventMatchingFilter =
9698
Object.entries(filter)
9799
.filter(([key, criteria]) => isGenericTagQuery(key) && Array.isArray(criteria))
98100
.some(([key, criteria]) => {
99-
return !event.tags.some((tag) => criteria.some((criterion) => isMatchingGenericTagCriterion(key, criterion)(tag)))
101+
return !event.tags.some((tag) =>
102+
criteria.some((criterion) => isMatchingGenericTagCriterion(key, criterion)(tag)),
103+
)
100104
})
101105
) {
102106
return false
@@ -205,6 +209,10 @@ export const isEphemeralEvent = (event: Event): boolean => {
205209
return event.kind >= EventKinds.EPHEMERAL_FIRST && event.kind <= EventKinds.EPHEMERAL_LAST
206210
}
207211

212+
export const isDvmJobRequestEvent = (event: Event): boolean => {
213+
return event.kind >= EventKinds.DVM_JOB_REQUEST_FIRST && event.kind <= EventKinds.DVM_JOB_REQUEST_LAST
214+
}
215+
208216
export const isParameterizedReplaceableEvent = (event: Event): boolean => {
209217
return (
210218
event.kind >= EventKinds.PARAMETERIZED_REPLACEABLE_FIRST && event.kind <= EventKinds.PARAMETERIZED_REPLACEABLE_LAST

‎test/unit/factories/event-strategy-factory.spec.ts‎

Lines changed: 28 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,14 @@
11
import { expect } from 'chai'
22

3-
import { IEventRepository, IInviteCodeRepository, IUserRepository } from '../../../src/@types/repositories'
3+
import {
4+
IDvmJobRepository,
5+
IEventRepository,
6+
IInviteCodeRepository,
7+
IUserRepository,
8+
} from '../../../src/@types/repositories'
49
import { DefaultEventStrategy } from '../../../src/handlers/event-strategies/default-event-strategy'
510
import { DeleteEventStrategy } from '../../../src/handlers/event-strategies/delete-event-strategy'
11+
import { DvmJobRequestEventStrategy } from '../../../src/handlers/event-strategies/dvm-job-request-event-strategy'
612
import { EphemeralEventStrategy } from '../../../src/handlers/event-strategies/ephemeral-event-strategy'
713
import { Event } from '../../../src/@types/event'
814
import { EventKinds } from '../../../src/constants/base'
@@ -24,6 +30,7 @@ describe('eventStrategyFactory', () => {
2430
let eventRepository: IEventRepository
2531
let userRepository: IUserRepository
2632
let inviteCodeRepository: IInviteCodeRepository
33+
let dvmJobRepository: IDvmJobRepository
2734
let cache: ICacheAdapter
2835
let settings: () => Settings
2936
let event: Event
@@ -34,12 +41,20 @@ describe('eventStrategyFactory', () => {
3441
eventRepository = {} as any
3542
userRepository = {} as any
3643
inviteCodeRepository = {} as any
44+
dvmJobRepository = {} as any
3745
cache = {} as any
38-
settings = () => ({ info: { relay_url: 'wss://test.relay' } } as any)
46+
settings = () => ({ info: { relay_url: 'wss://test.relay' } }) as any
3947
event = {} as any
4048
adapter = {} as any
4149

42-
factory = eventStrategyFactory(eventRepository, userRepository, inviteCodeRepository, cache, settings)
50+
factory = eventStrategyFactory(
51+
eventRepository,
52+
userRepository,
53+
inviteCodeRepository,
54+
dvmJobRepository,
55+
cache,
56+
settings,
57+
)
4358
})
4459

4560
it('returns ReplaceableEvent given a set_metadata event', () => {
@@ -136,4 +151,14 @@ describe('eventStrategyFactory', () => {
136151
event.kind = EventKinds.NIP43_LEAVE_REQUEST
137152
expect(factory([event, adapter])).to.be.an.instanceOf(LeaveRequestEventStrategy)
138153
})
154+
155+
it('returns DvmJobRequestEventStrategy given a DVM job request (kind 5000-5999)', () => {
156+
event.kind = EventKinds.DVM_JOB_REQUEST_FIRST
157+
expect(factory([event, adapter])).to.be.an.instanceOf(DvmJobRequestEventStrategy)
158+
})
159+
160+
it('returns DvmJobRequestEventStrategy given the last DVM job request kind (5999)', () => {
161+
event.kind = EventKinds.DVM_JOB_REQUEST_LAST
162+
expect(factory([event, adapter])).to.be.an.instanceOf(DvmJobRequestEventStrategy)
163+
})
139164
})

0 commit comments

Comments
 (0)