Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/dvm-job-ingestion.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"nostream": minor
---

feat(dvm): trap NIP-90 job request events (kind 5000-5999) and record them via the job repository
5 changes: 5 additions & 0 deletions .changeset/dvm-job-persistence.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"nostream": minor
---

feat(dvm): add job persistence migration and repository for DVM job state
23 changes: 23 additions & 0 deletions migrations/20260812_150000_create_dvm_jobs_table.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
exports.up = function (knex) {
return knex.schema.createTable('dvm_jobs', (table) => {
table.binary('id').primary()
table.binary('requester_pubkey').notNullable()
table.integer('kind').unsigned().notNullable()
table.integer('worker_index').nullable()
table.enum('status', ['submitted', 'picked_up', 'completed', 'failed', 'timed_out']).notNullable().defaultTo('submitted')
table.binary('result_event_id').nullable()
table.text('error').nullable()
table.timestamp('picked_up_at', { useTz: true }).nullable()
table.timestamp('completed_at', { useTz: true }).nullable()
table.timestamp('created_at', { useTz: true }).notNullable().defaultTo(knex.fn.now())
table.timestamp('updated_at', { useTz: true }).notNullable().defaultTo(knex.fn.now())

table.index(['requester_pubkey'], 'idx_dvm_jobs_requester_pubkey')
table.index(['status'], 'idx_dvm_jobs_status')
table.index(['kind'], 'idx_dvm_jobs_kind')
})
}

exports.down = function (knex) {
return knex.schema.dropTable('dvm_jobs')
}
37 changes: 37 additions & 0 deletions src/@types/dvm.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
import { Pubkey } from './base'

export enum DvmJobStatus {
SUBMITTED = 'submitted',
PICKED_UP = 'picked_up',
COMPLETED = 'completed',
FAILED = 'failed',
TIMED_OUT = 'timed_out',
}

export interface DvmJob {
id: string
requesterPubkey: Pubkey
kind: number
workerIndex: number | null
status: DvmJobStatus
resultEventId: string | null
error: string | null
pickedUpAt: Date | null
completedAt: Date | null
createdAt: Date
updatedAt: Date
}

export interface DBDvmJob {
id: Buffer
requester_pubkey: Buffer
kind: number
worker_index: number | null
status: DvmJobStatus
result_event_id: Buffer | null
error: string | null
picked_up_at: Date | null
completed_at: Date | null
created_at: Date
updated_at: Date
}
16 changes: 13 additions & 3 deletions src/@types/repositories.ts
Original file line number Diff line number Diff line change
@@ -1,12 +1,12 @@
import { PassThrough } from 'stream'

import { EventKinds } from '../constants/base'
import { DatabaseClient, EventId, Pubkey } from './base'
import { DvmJob } from './dvm'
import { DBEvent, Event } from './event'
import { EventKinds } from '../constants/base'
import { EventKindsRange } from './settings'
import { InviteCode } from './invite-code'
import { Invoice } from './invoice'
import { Nip05Verification } from './nip05'
import { EventKindsRange } from './settings'
import { SubscriptionFilter } from './subscription'
import { User } from './user'

Expand Down Expand Up @@ -73,3 +73,13 @@ export interface IInviteCodeRepository {
findActiveCodes(limit?: number): Promise<InviteCode[]>
deleteExpiredCodes(): Promise<number>
}

export interface IDvmJobRepository {
create(id: string, requesterPubkey: Pubkey, kind: number): Promise<DvmJob>
findById(id: string): Promise<DvmJob | undefined>
assignWorker(id: string, workerIndex: number): Promise<boolean>
updateStatus(
job: Pick<DvmJob, 'id' | 'status'> & Partial<Pick<DvmJob, 'resultEventId' | 'error'>>,
): Promise<DvmJob | undefined>
findPendingJobs(limit?: number): Promise<DvmJob[]>
}
3 changes: 3 additions & 0 deletions src/constants/base.ts
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,9 @@ export enum EventKinds {
// Lightning zaps
ZAP_REQUEST = 9734,
ZAP_RECEIPT = 9735,
// NIP-90: Data Vending Machines — job request events
DVM_JOB_REQUEST_FIRST = 5000,
DVM_JOB_REQUEST_LAST = 5999,
// Replaceable events
REPLACEABLE_FIRST = 10000,
// NIP-65: Relay List Metadata
Expand Down
16 changes: 12 additions & 4 deletions src/factories/event-strategy-factory.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,8 @@
import { ICacheAdapter, IWebSocketAdapter } from '../@types/adapters'
import { IEventRepository, IInviteCodeRepository, IUserRepository } from '../@types/repositories'
import { IDvmJobRepository, IEventRepository, IInviteCodeRepository, IUserRepository } from '../@types/repositories'
import {
isDeleteEvent,
isDvmJobRequestEvent,
isEphemeralEvent,
isGiftWrapEvent,
isMarmotGroupEvent,
Expand All @@ -14,6 +15,7 @@ import { isNip43JoinRequest, isNip43LeaveRequest } from '../utils/nip43'
import { isRelayListEvent } from '../utils/nip65'
import { DefaultEventStrategy } from '../handlers/event-strategies/default-event-strategy'
import { DeleteEventStrategy } from '../handlers/event-strategies/delete-event-strategy'
import { DvmJobRequestEventStrategy } from '../handlers/event-strategies/dvm-job-request-event-strategy'
import { EphemeralEventStrategy } from '../handlers/event-strategies/ephemeral-event-strategy'
import { Event } from '../@types/event'
import { Factory } from '../@types/base'
Expand All @@ -33,6 +35,7 @@ export const eventStrategyFactory =
eventRepository: IEventRepository,
userRepository: IUserRepository,
inviteCodeRepository: IInviteCodeRepository,
dvmJobRepository: IDvmJobRepository,
cache: ICacheAdapter,
settings: () => Settings,
): Factory<IEventStrategy<Event, Promise<void>>, [Event, IWebSocketAdapter]> =>
Expand All @@ -47,12 +50,17 @@ export const eventStrategyFactory =
return new TimestampEventStrategy(adapter, eventRepository)
} else if (isRelayListEvent(event) || isReplaceableEvent(event)) {
return new ReplaceableEventStrategy(adapter, eventRepository)
// NIP-43: Join/Leave requests MUST be checked before the generic ephemeral
// handler, because kinds 28934/28936 fall in the ephemeral range (20000-29999).
// NIP-43: Join/Leave requests MUST be checked before the generic ephemeral
// handler, because kinds 28934/28936 fall in the ephemeral range (20000-29999).
} else if (isNip43JoinRequest(event)) {
return new JoinRequestEventStrategy(adapter, inviteCodeRepository, userRepository, cache, settings)
} else if (isNip43LeaveRequest(event)) {
return new LeaveRequestEventStrategy(adapter, userRepository, cache, settings)
// NIP-90: DVM job requests (kind 5000-5999) checked early, same reasoning
// as the NIP-43 checks above — kept explicit rather than relying on it
// falling through to DefaultEventStrategy.
} else if (isDvmJobRequestEvent(event)) {
return new DvmJobRequestEventStrategy(adapter, eventRepository, dvmJobRepository)
} else if (isEphemeralEvent(event)) {
return new EphemeralEventStrategy(adapter)
} else if (isDeleteEvent(event)) {
Expand All @@ -62,4 +70,4 @@ export const eventStrategyFactory =
}

return new DefaultEventStrategy(adapter, eventRepository)
}
}
18 changes: 16 additions & 2 deletions src/factories/message-handler-factory.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,11 @@
import { ICacheAdapter, IWebSocketAdapter } from '../@types/adapters'
import { IEventRepository, IInviteCodeRepository, INip05VerificationRepository, IUserRepository } from '../@types/repositories'
import {
IDvmJobRepository,
IEventRepository,
IInviteCodeRepository,
INip05VerificationRepository,
IUserRepository,
} from '../@types/repositories'
import { IncomingMessage, MessageType } from '../@types/messages'
import { createSettings } from './settings-factory'
import { AuthMessageHandler } from '../handlers/auth-message-handler'
Expand All @@ -26,13 +32,21 @@ export const messageHandlerFactory =
userRepository: IUserRepository,
nip05VerificationRepository: INip05VerificationRepository,
inviteCodeRepository: IInviteCodeRepository,
dvmJobRepository: IDvmJobRepository,
) =>
([message, adapter]: [IncomingMessage, IWebSocketAdapter]) => {
switch (message[0]) {
case MessageType.EVENT: {
return new EventMessageHandler(
adapter,
eventStrategyFactory(eventRepository, userRepository, inviteCodeRepository, getCache(), createSettings),
eventStrategyFactory(
eventRepository,
userRepository,
inviteCodeRepository,
dvmJobRepository,
getCache(),
createSettings,
),
eventRepository,
userRepository,
createSettings,
Expand Down
17 changes: 15 additions & 2 deletions src/factories/websocket-adapter-factory.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,13 @@
import { IncomingMessage } from 'http'
import { WebSocket } from 'ws'

import { IEventRepository, IInviteCodeRepository, INip05VerificationRepository, IUserRepository } from '../@types/repositories'
import {
IDvmJobRepository,
IEventRepository,
IInviteCodeRepository,
INip05VerificationRepository,
IUserRepository,
} from '../@types/repositories'
import { createSettings } from './settings-factory'
import { IWebSocketServerAdapter } from '../@types/adapters'
import { messageHandlerFactory } from './message-handler-factory'
Expand All @@ -14,13 +20,20 @@ export const webSocketAdapterFactory =
userRepository: IUserRepository,
nip05VerificationRepository: INip05VerificationRepository,
inviteCodeRepository: IInviteCodeRepository,
dvmJobRepository: IDvmJobRepository,
) =>
([client, request, webSocketServerAdapter]: [WebSocket, IncomingMessage, IWebSocketServerAdapter]) =>
new WebSocketAdapter(
client,
request,
webSocketServerAdapter,
messageHandlerFactory(eventRepository, userRepository, nip05VerificationRepository, inviteCodeRepository),
messageHandlerFactory(
eventRepository,
userRepository,
nip05VerificationRepository,
inviteCodeRepository,
dvmJobRepository,
),
rateLimiterFactory,
createSettings,
)
10 changes: 9 additions & 1 deletion src/factories/worker-factory.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import { AppWorker } from '../app/worker'
import { createLogger } from './logger-factory'
import { createSettings } from '../factories/settings-factory'
import { createWebApp } from './web-app-factory'
import { DvmJobRepository } from '../repositories/dvm-job-repository'
import { EventRepository } from '../repositories/event-repository'
import { InviteCodeRepository } from '../repositories/invite-code-repository'
import { Nip05VerificationRepository } from '../repositories/nip05-verification-repository'
Expand All @@ -24,6 +25,7 @@ export const workerFactory = (): AppWorker => {
const userRepository = new UserRepository(dbClient, eventRepository)
const nip05VerificationRepository = new Nip05VerificationRepository(dbClient)
const inviteCodeRepository = new InviteCodeRepository(dbClient)
const dvmJobRepository = new DvmJobRepository(dbClient)

const settings = createSettings()

Expand Down Expand Up @@ -65,7 +67,13 @@ export const workerFactory = (): AppWorker => {
const adapter = new WebSocketServerAdapter(
server,
webSocketServer,
webSocketAdapterFactory(eventRepository, userRepository, nip05VerificationRepository, inviteCodeRepository),
webSocketAdapterFactory(
eventRepository,
userRepository,
nip05VerificationRepository,
inviteCodeRepository,
dvmJobRepository,
),
createSettings,
)

Expand Down
42 changes: 42 additions & 0 deletions src/handlers/event-strategies/dvm-job-request-event-strategy.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
import { createEventCommandResult } from '../../telemetry/event-metrics'
import { createLogger } from '../../factories/logger-factory'
import { Event } from '../../@types/event'
import { IDvmJobRepository, IEventRepository } from '../../@types/repositories'
import { IEventStrategy } from '../../@types/message-handlers'
import { IWebSocketAdapter } from '../../@types/adapters'
import { WebSocketAdapterEvent } from '../../constants/adapter'

const logger = createLogger('dvm-job-request-event-strategy')

export class DvmJobRequestEventStrategy implements IEventStrategy<Event, Promise<void>> {
public constructor(
private readonly webSocket: IWebSocketAdapter,
private readonly eventRepository: IEventRepository,
private readonly dvmJobRepository: IDvmJobRepository,
) {}

public async execute(event: Event): Promise<void> {
logger('received dvm job request: %o', event)

const count = await this.eventRepository.create(event)
this.webSocket.emit(
WebSocketAdapterEvent.Message,
createEventCommandResult(event.id, true, count ? '' : 'duplicate:'),
)

if (!count) {
return
}

this.webSocket.emit(WebSocketAdapterEvent.Broadcast, event)

try {
await this.dvmJobRepository.create(event.id, event.pubkey, event.kind)
} catch (error) {
// Job-state recording is best-effort: the event itself is already
// stored and broadcast correctly, so a repository failure here must
// not surface as a rejection of a valid event.
logger.error('unable to record dvm job for event %s: %o', event.id, error)
}
}
}
Loading
Loading