From ca8c8708ee848aeeb49dd9dfc0fb5bc0ffc0f2df Mon Sep 17 00:00:00 2001 From: Priyanshubhartistm Date: Thu, 6 Aug 2026 22:30:02 +0530 Subject: [PATCH 1/4] feat(dvm): add dvm worker registry settings Signed-off-by: Priyanshubhartistm --- resources/default-settings.yaml | 2 ++ src/@types/settings.ts | 16 ++++++++++++++++ 2 files changed, 18 insertions(+) diff --git a/resources/default-settings.yaml b/resources/default-settings.yaml index 40e72178..9eba3b9f 100755 --- a/resources/default-settings.yaml +++ b/resources/default-settings.yaml @@ -107,6 +107,8 @@ workers: count: 0 mirroring: static: [] +dvm: + workers: [] limits: # strategy selection configuration for rate limiting: rateLimiter: diff --git a/src/@types/settings.ts b/src/@types/settings.ts index 830a5dbe..05fa0a8d 100644 --- a/src/@types/settings.ts +++ b/src/@types/settings.ts @@ -246,6 +246,21 @@ export interface Mirroring { static?: Mirror[] } +export interface DvmWorker { + /** Command to spawn for this worker (e.g. an interpreter or executable path). */ + command: string + /** Arguments passed to the spawned command. */ + args?: string[] + /** NIP-90 job request kinds (5000-5999) this worker accepts. */ + kinds?: number[] + /** Max time in ms to wait for a job result before considering it timed out. */ + timeoutMs?: number +} + +export interface Dvm { + workers?: DvmWorker[] +} + export type Nip05Mode = 'enabled' | 'passive' | 'disabled' export interface Nip45Settings { @@ -330,6 +345,7 @@ export interface Settings { workers?: Worker limits?: Limits mirroring?: Mirroring + dvm?: Dvm nip05?: Nip05Settings nip42?: Nip42Settings nip43?: Nip43Settings From be4577b6ceccfcac417840cb4fdaf0a47b370995 Mon Sep 17 00:00:00 2001 From: Priyanshubhartistm Date: Thu, 6 Aug 2026 22:30:49 +0530 Subject: [PATCH 2/4] feat(dvm): wire WORKER_TYPE=dvm-orchestrator process topology Signed-off-by: Priyanshubhartistm --- src/app/app.ts | 12 ++++ src/app/dvm-orchestrator-worker.ts | 58 +++++++++++++++++++ .../dvm-orchestrator-worker-factory.ts | 7 +++ src/index.ts | 5 +- 4 files changed, 81 insertions(+), 1 deletion(-) create mode 100644 src/app/dvm-orchestrator-worker.ts create mode 100644 src/factories/dvm-orchestrator-worker-factory.ts diff --git a/src/app/app.ts b/src/app/app.ts index cccdc607..cb247902 100644 --- a/src/app/app.ts +++ b/src/app/app.ts @@ -108,6 +108,18 @@ export class App implements IRunnable { logCentered(`${mirrors.length} static-mirroring worker started`, width) } + const dvmWorkers = settings?.dvm?.workers + + if (Array.isArray(dvmWorkers) && dvmWorkers.length) { + for (let i = 0; i < dvmWorkers.length; i++) { + createWorker({ + WORKER_TYPE: 'dvm-orchestrator', + DVM_WORKER_INDEX: i.toString(), + }) + } + logCentered(`${dvmWorkers.length} dvm-orchestrator worker started`, width) + } + logger('settings: %O', settings) const host = `${hostname()}:${port}` diff --git a/src/app/dvm-orchestrator-worker.ts b/src/app/dvm-orchestrator-worker.ts new file mode 100644 index 00000000..4947f722 --- /dev/null +++ b/src/app/dvm-orchestrator-worker.ts @@ -0,0 +1,58 @@ +import { path } from 'ramda' +import { IRunnable } from '../@types/base' +import { DvmWorker, Settings } from '../@types/settings' +import { createLogger } from '../factories/logger-factory' +import { shutdownMetricsTelemetry } from '../telemetry/metrics' + +const logger = createLogger('dvm-orchestrator-worker') + +export class DvmOrchestratorWorker implements IRunnable { + private config: DvmWorker | undefined + + public constructor( + private readonly process: NodeJS.Process, + private readonly settings: () => Settings, + ) { + this.process + .on('SIGINT', this.onExit.bind(this)) + .on('SIGHUP', this.onExit.bind(this)) + .on('SIGTERM', this.onExit.bind(this)) + .on('uncaughtException', this.onError.bind(this)) + .on('unhandledRejection', this.onError.bind(this)) + } + + public run(): void { + const currentSettings = this.settings() + + this.config = path(['dvm', 'workers', this.process.env.DVM_WORKER_INDEX], currentSettings) as DvmWorker | undefined + + if (!this.config) { + logger.error('no dvm worker config found for index %s', this.process.env.DVM_WORKER_INDEX) + this.process.exit(1) + return + } + + logger.info('dvm-orchestrator worker started for command: %s', this.config.command) + } + + private onError(error: Error) { + logger('error: %o', error) + throw error + } + + private onExit() { + logger('exiting') + void shutdownMetricsTelemetry().finally(() => { + this.close(() => { + this.process.exit(0) + }) + }) + } + + public close(callback?: () => void) { + logger('closing') + if (typeof callback === 'function') { + callback() + } + } +} diff --git a/src/factories/dvm-orchestrator-worker-factory.ts b/src/factories/dvm-orchestrator-worker-factory.ts new file mode 100644 index 00000000..47657d9d --- /dev/null +++ b/src/factories/dvm-orchestrator-worker-factory.ts @@ -0,0 +1,7 @@ +import process from 'process' +import { DvmOrchestratorWorker } from '../app/dvm-orchestrator-worker' +import { createSettings } from './settings-factory' + +export const dvmOrchestratorWorkerFactory = () => { + return new DvmOrchestratorWorker(process, createSettings) +} diff --git a/src/index.ts b/src/index.ts index cdae9e1e..d43c56a3 100644 --- a/src/index.ts +++ b/src/index.ts @@ -1,10 +1,11 @@ import cluster from 'cluster' import { appFactory } from './factories/app-factory' +import { dvmOrchestratorWorkerFactory } from './factories/dvm-orchestrator-worker-factory' import { maintenanceWorkerFactory } from './factories/maintenance-worker-factory' import { staticMirroringWorkerFactory } from './factories/static-mirroring.worker-factory' -import { initializeMetricsTelemetry } from './telemetry/metrics' import { workerFactory } from './factories/worker-factory' +import { initializeMetricsTelemetry } from './telemetry/metrics' export const getRunner = () => { if (cluster.isPrimary) { @@ -17,6 +18,8 @@ export const getRunner = () => { return maintenanceWorkerFactory() case 'static-mirroring': return staticMirroringWorkerFactory() + case 'dvm-orchestrator': + return dvmOrchestratorWorkerFactory() default: throw new Error(`Unknown worker: ${process.env.WORKER_TYPE}`) } From b8542dfae4dc447e5cb11eb1925855d5f1c5abf9 Mon Sep 17 00:00:00 2001 From: Priyanshubhartistm Date: Thu, 6 Aug 2026 22:31:36 +0530 Subject: [PATCH 3/4] test(dvm): add dvm-orchestrator-worker unit tests Signed-off-by: Priyanshubhartistm --- test/unit/app/dvm-orchestrator-worker.spec.ts | 87 +++++++++++++++++++ .../dvm-orchestrator-worker-factory.spec.ts | 10 +++ 2 files changed, 97 insertions(+) create mode 100644 test/unit/app/dvm-orchestrator-worker.spec.ts create mode 100644 test/unit/factories/dvm-orchestrator-worker-factory.spec.ts diff --git a/test/unit/app/dvm-orchestrator-worker.spec.ts b/test/unit/app/dvm-orchestrator-worker.spec.ts new file mode 100644 index 00000000..f53cc9cd --- /dev/null +++ b/test/unit/app/dvm-orchestrator-worker.spec.ts @@ -0,0 +1,87 @@ +import chai from 'chai' +import EventEmitter from 'events' +import Sinon from 'sinon' +import sinonChai from 'sinon-chai' + +import { Settings } from '../../../src/@types/settings' +import { DvmOrchestratorWorker } from '../../../src/app/dvm-orchestrator-worker' +import * as metricsTelemetry from '../../../src/telemetry/metrics' + +chai.use(sinonChai) + +const { expect } = chai + +describe('DvmOrchestratorWorker', () => { + let sandbox: Sinon.SinonSandbox + let fakeProcess: EventEmitter & { exit: Sinon.SinonStub; env: Record } + let settings: Sinon.SinonStub + let settingsState: Settings + + beforeEach(() => { + sandbox = Sinon.createSandbox() + + fakeProcess = Object.assign(new EventEmitter(), { + exit: sandbox.stub(), + env: {}, + }) as EventEmitter & { exit: Sinon.SinonStub; env: Record } + + settingsState = { + dvm: { + workers: [{ command: 'python3', args: ['worker.py'] }], + }, + } as any + + settings = sandbox.stub().callsFake(() => settingsState) + + sandbox.stub(metricsTelemetry, 'shutdownMetricsTelemetry').resolves() + }) + + afterEach(() => { + sandbox.restore() + }) + + describe('run', () => { + it('logs startup for the worker config at DVM_WORKER_INDEX', () => { + fakeProcess.env.DVM_WORKER_INDEX = '0' + const worker = new DvmOrchestratorWorker(fakeProcess as any, settings as any) + + expect(() => worker.run()).to.not.throw() + expect(fakeProcess.exit).not.to.have.been.called + }) + + it('exits with code 1 if no worker config exists for the given index', () => { + fakeProcess.env.DVM_WORKER_INDEX = '5' + const worker = new DvmOrchestratorWorker(fakeProcess as any, settings as any) + + worker.run() + + expect(fakeProcess.exit).to.have.been.calledWith(1) + }) + }) + + describe('signal handling', () => { + it('closes and exits on SIGTERM', async () => { + fakeProcess.env.DVM_WORKER_INDEX = '0' + const worker = new DvmOrchestratorWorker(fakeProcess as any, settings as any) + worker.run() + + fakeProcess.emit('SIGTERM') + + await Promise.resolve() + await Promise.resolve() + + expect(fakeProcess.exit).to.have.been.calledWith(0) + }) + }) + + describe('close', () => { + it('invokes the callback', () => { + const worker = new DvmOrchestratorWorker(fakeProcess as any, settings as any) + const callback = sandbox.stub() + + worker.close(callback) + + expect(callback).to.have.been.calledOnce + }) + }) +}) diff --git a/test/unit/factories/dvm-orchestrator-worker-factory.spec.ts b/test/unit/factories/dvm-orchestrator-worker-factory.spec.ts new file mode 100644 index 00000000..babbdc58 --- /dev/null +++ b/test/unit/factories/dvm-orchestrator-worker-factory.spec.ts @@ -0,0 +1,10 @@ +import { expect } from 'chai' + +import { DvmOrchestratorWorker } from '../../../src/app/dvm-orchestrator-worker' +import { dvmOrchestratorWorkerFactory } from '../../../src/factories/dvm-orchestrator-worker-factory' + +describe('dvmOrchestratorWorkerFactory', () => { + it('returns a DvmOrchestratorWorker', () => { + expect(dvmOrchestratorWorkerFactory()).to.be.an.instanceOf(DvmOrchestratorWorker) + }) +}) From 8257e7feb4929916ce52d784a390c312af791461 Mon Sep 17 00:00:00 2001 From: Priyanshubhartistm Date: Thu, 6 Aug 2026 22:32:31 +0530 Subject: [PATCH 4/4] docs: document dvm.workers settings Signed-off-by: Priyanshubhartistm --- .changeset/dvm-worker-registry.md | 5 +++++ CONFIGURATION.md | 4 ++++ 2 files changed, 9 insertions(+) create mode 100644 .changeset/dvm-worker-registry.md diff --git a/.changeset/dvm-worker-registry.md b/.changeset/dvm-worker-registry.md new file mode 100644 index 00000000..31e4284e --- /dev/null +++ b/.changeset/dvm-worker-registry.md @@ -0,0 +1,5 @@ +--- +"nostream": minor +--- + +feat(dvm): add worker registry settings and dvm-orchestrator process topology diff --git a/CONFIGURATION.md b/CONFIGURATION.md index 4604edcc..3e5ff311 100644 --- a/CONFIGURATION.md +++ b/CONFIGURATION.md @@ -133,6 +133,10 @@ The settings below are listed in alphabetical order by name. Please keep this ta | Name | Description | |---------------------------------------------|-------------------------------------------------------------------------------| +| dvm.workers[].args | Arguments passed to the spawned command. Optional. | +| dvm.workers[].command | Command to spawn for this DVM worker (e.g. an interpreter or executable path). | +| dvm.workers[].kinds | NIP-90 job request kinds (5000-5999) this worker accepts. Optional. | +| dvm.workers[].timeoutMs | Max time in ms to wait for a job result before considering it timed out. Optional. | | info.banner | Public banner image URL for the relay information document. | | info.contact | Relay operator's contact. (e.g. mailto:operator@relay-your-domain.com) | | info.description | Public description of your relay. (e.g. Toronto Bitcoin Group Public Relay) |