diff --git a/.changeset/dvm-job-dispatch.md b/.changeset/dvm-job-dispatch.md new file mode 100644 index 00000000..7ffab62f --- /dev/null +++ b/.changeset/dvm-job-dispatch.md @@ -0,0 +1,5 @@ +--- +"nostream": minor +--- + +feat(dvm): dispatch pending DVM jobs to worker processes and publish kind 6000-6999 results back diff --git a/src/@types/repositories.ts b/src/@types/repositories.ts index c25a0e62..d4c82a1a 100644 --- a/src/@types/repositories.ts +++ b/src/@types/repositories.ts @@ -81,5 +81,5 @@ export interface IDvmJobRepository { updateStatus( job: Pick & Partial>, ): Promise - findPendingJobs(limit?: number): Promise + findPendingJobs(limit?: number, kinds?: number[]): Promise } diff --git a/src/app/dvm-orchestrator-worker.ts b/src/app/dvm-orchestrator-worker.ts index 4947f722..358de3f9 100644 --- a/src/app/dvm-orchestrator-worker.ts +++ b/src/app/dvm-orchestrator-worker.ts @@ -1,17 +1,48 @@ -import { path } from 'ramda' -import { IRunnable } from '../@types/base' +import { andThen, otherwise, path, pipe } from 'ramda' + +import { + broadcastEvent, + getPublicKey, + getRelayPrivateKey, + identifyEvent, + signEvent, + toNostrEvent, +} from '../utils/event' +import { DvmJob, DvmJobStatus } from '../@types/dvm' import { DvmWorker, Settings } from '../@types/settings' +import { Event, UnidentifiedEvent } from '../@types/event' +import { EventKinds, EventTags } from '../constants/base' +import { IDvmJobRepository, IEventRepository } from '../@types/repositories' +import { spawnWorkerProcess, WorkerProcessHandle, WorkerSpawnErrorReason } from '../cli/utils/process' import { createLogger } from '../factories/logger-factory' +import { IRunnable } from '../@types/base' import { shutdownMetricsTelemetry } from '../telemetry/metrics' const logger = createLogger('dvm-orchestrator-worker') +const POLL_INTERVAL_MS = 2000 +const DEFAULT_JOB_TIMEOUT_MS = 60000 +const DISPATCH_BATCH_SIZE = 10 + +type PendingJob = { + job: DvmJob + requestEvent: Event + timer: NodeJS.Timeout +} + export class DvmOrchestratorWorker implements IRunnable { private config: DvmWorker | undefined + private interval: NodeJS.Timeout | undefined + private isRunning = false + private closing = false + private worker: WorkerProcessHandle | undefined + private readonly pending = new Map() public constructor( private readonly process: NodeJS.Process, private readonly settings: () => Settings, + private readonly dvmJobRepository: IDvmJobRepository, + private readonly eventRepository: IEventRepository, ) { this.process .on('SIGINT', this.onExit.bind(this)) @@ -33,6 +64,221 @@ export class DvmOrchestratorWorker implements IRunnable { } logger.info('dvm-orchestrator worker started for command: %s', this.config.command) + + this.ensureWorkerProcess() + + this.interval = setInterval(async () => { + if (this.isRunning) { + logger('skipping scheduled dispatch because previous run is still in progress') + return + } + + this.isRunning = true + try { + await this.dispatchNextJob() + } catch (error) { + this.onError(error as Error) + } finally { + this.isRunning = false + } + }, POLL_INTERVAL_MS) + } + + // One long-lived worker process per configured dvm.workers[i], multiplexing + // every in-flight job over a single newline-delimited-JSON stdin/stdout pipe + // (per issue #731 — process.ts's one-shot spawn helpers buffer all output + // into a single string, which doesn't work once more than one job can be + // in flight against the same worker at a time). + private ensureWorkerProcess(): void { + if (this.worker || !this.config || this.closing) { + return + } + + const worker = spawnWorkerProcess(this.config.command, this.config.args ?? []) + worker.onMessage((message) => this.handleWorkerMessage(message)) + worker.onExit((code, signal) => this.handleWorkerExit(code, signal)) + worker.onSpawnError((reason) => this.handleWorkerSpawnError(reason)) + this.worker = worker + } + + private async dispatchNextJob(): Promise { + if (!this.config) { + return + } + + // Lazily respawn if the worker process died since the last tick. + this.ensureWorkerProcess() + if (!this.worker) { + return + } + + // findPendingJobs() returns both SUBMITTED and PICKED_UP jobs (oldest first); + // assignWorker() only succeeds against SUBMITTED ones. Walk a batch instead of + // just the oldest candidate so an already-picked-up job at the head can't + // starve out later SUBMITTED jobs behind it. + const candidates = await this.dvmJobRepository.findPendingJobs(DISPATCH_BATCH_SIZE, this.config.kinds) + + let job: DvmJob | undefined + for (const candidate of candidates) { + if (await this.dvmJobRepository.assignWorker(candidate.id, this.workerIndex())) { + job = candidate + break + } + } + + if (!job) { + // Lost the race on every candidate in this batch — try again next tick. + return + } + + logger('picked up job %s (kind %d)', job.id, job.kind) + + const [row] = await this.eventRepository.findByFilters([{ ids: [job.id] }]) + if (!row) { + await this.failJob(job.id, 'source event not found') + return + } + + const requestEvent = toNostrEvent(row) + const timeoutMs = this.config.timeoutMs ?? DEFAULT_JOB_TIMEOUT_MS + const timer = setTimeout(() => this.handleJobTimeout(job.id), timeoutMs) + this.pending.set(job.id, { job, requestEvent, timer }) + + const sent = this.worker.send({ + id: requestEvent.id, + kind: requestEvent.kind, + pubkey: requestEvent.pubkey, + tags: requestEvent.tags, + content: requestEvent.content, + }) + + if (!sent) { + clearTimeout(timer) + this.pending.delete(job.id) + await this.failJob(job.id, 'unable to send job to worker process') + } + } + + // The worker replies on the same pipe with { id: , content: }, + // correlated back to the job that's still pending — replies for jobs we no + // longer track (already timed out, or from a previous worker instance) are + // dropped rather than treated as an error. + private handleWorkerMessage(message: unknown): void { + const jobId = (message as { id?: unknown } | null)?.id + if (typeof jobId !== 'string') { + logger.error('ignoring malformed worker message (missing id): %o', message) + return + } + + const pending = this.pending.get(jobId) + if (!pending) { + return + } + + clearTimeout(pending.timer) + this.pending.delete(jobId) + + const rawContent = (message as { content?: unknown }).content + const content = typeof rawContent === 'string' ? rawContent : JSON.stringify(message) + + void this.publishResult(pending.job, pending.requestEvent, content) + } + + private handleJobTimeout(jobId: string): void { + const pending = this.pending.get(jobId) + if (!pending) { + return + } + + this.pending.delete(jobId) + void this.failJob(jobId, 'worker timeout', true) + + // Soft Node-level guard per issue #731: kill the worker process on timeout. + // Any other jobs still in flight on it fail as a side effect of the exit + // handler below; ensureWorkerProcess() respawns on the next dispatch tick. + this.worker?.kill() + } + + private handleWorkerExit(code: number | null, signal: NodeJS.Signals | null): void { + logger.error('dvm worker process exited (code=%s, signal=%s)', code, signal) + this.worker = undefined + + for (const [jobId, pending] of this.pending) { + clearTimeout(pending.timer) + void this.failJob(jobId, `worker exited (code=${code ?? 'null'}, signal=${signal ?? 'null'})`) + } + this.pending.clear() + } + + private handleWorkerSpawnError(reason: WorkerSpawnErrorReason): void { + logger.error('unable to spawn dvm worker process: %s', reason) + this.worker = undefined + } + + private workerIndex(): number { + return Number(this.process.env.DVM_WORKER_INDEX) + } + + private async failJob(jobId: string, error: string, timedOut = false): Promise { + logger.error('job %s failed: %s', jobId, error) + try { + await this.dvmJobRepository.updateStatus({ + id: jobId, + status: timedOut ? DvmJobStatus.TIMED_OUT : DvmJobStatus.FAILED, + error, + }) + } catch (updateError) { + logger.error('unable to update failed job %s: %o', jobId, updateError) + } + } + + // Relay-authored, matching the same self-signing pattern used for invoice + // notifications (payments-service.ts) and NIP-89 authoring: the settings + // schema has no per-worker signing key, so the relay's own derived keypair + // is the DVM identity for every locally-bridged worker. + private async publishResult(job: DvmJob, requestEvent: Event, content: string): Promise { + const currentSettings = this.settings() + const relayPrivkey = getRelayPrivateKey(currentSettings.info.relay_url) + const relayPubkey = getPublicKey(relayPrivkey) + + const unsignedEvent: UnidentifiedEvent = { + pubkey: relayPubkey, + kind: (requestEvent.kind + 1000) as EventKinds, + created_at: Math.floor(Date.now() / 1000), + content, + tags: [ + [EventTags.Event, requestEvent.id], + [EventTags.Pubkey, requestEvent.pubkey], + ], + } + + const persistEvent = async (event: Event) => { + await this.eventRepository.create(event) + return event + } + + const markCompleted = async (event: Event) => { + await this.dvmJobRepository.updateStatus({ + id: job.id, + status: DvmJobStatus.COMPLETED, + resultEventId: event.id, + }) + return event + } + + const logPublishError = async (error: Error) => { + logger.error('unable to publish result for job %s: %o', job.id, error) + await this.failJob(job.id, `unable to publish result: ${error.message}`) + } + + await pipe( + identifyEvent, + andThen(signEvent(relayPrivkey)), + andThen(persistEvent), + andThen(broadcastEvent), + andThen(markCompleted), + otherwise(logPublishError), + )(unsignedEvent) } private onError(error: Error) { @@ -51,8 +297,33 @@ export class DvmOrchestratorWorker implements IRunnable { public close(callback?: () => void) { logger('closing') - if (typeof callback === 'function') { - callback() + this.closing = true + if (this.interval) { + clearInterval(this.interval) + } + + const invokeCallback = () => { + if (typeof callback === 'function') { + callback() + } + } + + if (this.pending.size === 0) { + this.worker?.kill() + invokeCallback() + return } + + // Fail every still-in-flight job before killing the worker: clearing + // `this.pending` first (as before) meant handleWorkerExit() found nothing + // to fail, leaving those jobs stuck at PICKED_UP in the DB forever. + const failures = Array.from(this.pending.entries()).map(([jobId, pending]) => { + clearTimeout(pending.timer) + return this.failJob(jobId, 'worker shutting down') + }) + this.pending.clear() + this.worker?.kill() + + void Promise.all(failures).finally(invokeCallback) } } diff --git a/src/cli/utils/process.ts b/src/cli/utils/process.ts index 51185144..47745ffe 100644 --- a/src/cli/utils/process.ts +++ b/src/cli/utils/process.ts @@ -9,7 +9,12 @@ export type RunOptions = { export type CommandResult = | { ok: true; code: number; stdout: string; stderr: string } - | { ok: false; reason: 'not-found' | 'permission-denied' | 'spawn-error' | 'timeout' | 'signal'; stdout: string; stderr: string } + | { + ok: false + reason: 'not-found' | 'permission-denied' | 'spawn-error' | 'timeout' | 'signal' + stdout: string + stderr: string + } export const runCommand = (command: string, args: string[], options: RunOptions = {}): Promise => { return new Promise((resolve, reject) => { @@ -80,7 +85,9 @@ export const runCommandWithOutput = ( }) child.on('error', (err: NodeJS.ErrnoException) => { - if (timer) { clearTimeout(timer) } + if (timer) { + clearTimeout(timer) + } if (err.code === 'ENOENT') { settle({ ok: false, reason: 'not-found', stdout, stderr }) } else if (err.code === 'EACCES') { @@ -91,7 +98,9 @@ export const runCommandWithOutput = ( }) child.on('close', (code, signal) => { - if (timer) { clearTimeout(timer) } + if (timer) { + clearTimeout(timer) + } if (timedOut) { settle({ ok: false, reason: 'timeout', stdout, stderr }) @@ -107,3 +116,91 @@ export const runCommandWithOutput = ( }) }) } + +export type WorkerSpawnErrorReason = 'not-found' | 'permission-denied' | 'spawn-error' + +export type WorkerProcessHandle = { + send: (message: unknown) => boolean + onMessage: (handler: (message: unknown) => void) => void + onExit: (handler: (code: number | null, signal: NodeJS.Signals | null) => void) => void + onSpawnError: (handler: (reason: WorkerSpawnErrorReason) => void) => void + kill: () => void +} + +/** + * Spawns a long-lived worker process and frames messages to/from it as + * newline-delimited JSON over stdin/stdout, so one process can multiplex + * multiple in-flight jobs. Reuses the same spawn/error-classification + * conventions as runCommandWithOutput rather than a second spawn pattern. + */ +export const spawnWorkerProcess = ( + command: string, + args: string[], + options: Pick = {}, +): WorkerProcessHandle => { + const child = spawn(command, args, { + cwd: options.cwd, + env: { ...process.env, ...(options.env ?? {}) }, + stdio: 'pipe', + shell: false, + }) + + const messageHandlers: Array<(message: unknown) => void> = [] + const exitHandlers: Array<(code: number | null, signal: NodeJS.Signals | null) => void> = [] + const spawnErrorHandlers: Array<(reason: WorkerSpawnErrorReason) => void> = [] + + let buffer = '' + child.stdout.on('data', (chunk) => { + buffer += chunk.toString() + let newlineIndex = buffer.indexOf('\n') + while (newlineIndex >= 0) { + const line = buffer.slice(0, newlineIndex).trim() + buffer = buffer.slice(newlineIndex + 1) + if (line) { + try { + const message = JSON.parse(line) + messageHandlers.forEach((handler) => handler(message)) + } catch { + // Malformed line from the worker — drop it rather than crash the relay process. + } + } + newlineIndex = buffer.indexOf('\n') + } + }) + + child.on('error', (err: NodeJS.ErrnoException) => { + const reason: WorkerSpawnErrorReason = + err.code === 'ENOENT' ? 'not-found' : err.code === 'EACCES' ? 'permission-denied' : 'spawn-error' + spawnErrorHandlers.forEach((handler) => handler(reason)) + }) + + child.on('exit', (code, signal) => { + exitHandlers.forEach((handler) => handler(code, signal)) + }) + + return { + send: (message) => { + // stdin.write() can throw synchronously (e.g. write after end) if the + // worker has already exited — report failure instead of crashing the + // orchestrator so the caller can fail the job. + try { + child.stdin.write(`${JSON.stringify(message)}\n`) + return true + } catch { + return false + } + }, + onMessage: (handler) => { + messageHandlers.push(handler) + }, + onExit: (handler) => { + exitHandlers.push(handler) + }, + onSpawnError: (handler) => { + spawnErrorHandlers.push(handler) + }, + kill: () => { + child.kill('SIGTERM') + }, + } +} diff --git a/src/constants/base.ts b/src/constants/base.ts index 0715ecc0..e23112cb 100644 --- a/src/constants/base.ts +++ b/src/constants/base.ts @@ -44,6 +44,9 @@ export enum EventKinds { // NIP-90: Data Vending Machines — job request events DVM_JOB_REQUEST_FIRST = 5000, DVM_JOB_REQUEST_LAST = 5999, + // NIP-90: Data Vending Machines — job result events + DVM_JOB_RESULT_FIRST = 6000, + DVM_JOB_RESULT_LAST = 6999, // Replaceable events REPLACEABLE_FIRST = 10000, // NIP-65: Relay List Metadata diff --git a/src/factories/dvm-orchestrator-worker-factory.ts b/src/factories/dvm-orchestrator-worker-factory.ts index 47657d9d..1ca5aa69 100644 --- a/src/factories/dvm-orchestrator-worker-factory.ts +++ b/src/factories/dvm-orchestrator-worker-factory.ts @@ -1,7 +1,16 @@ import process from 'process' -import { DvmOrchestratorWorker } from '../app/dvm-orchestrator-worker' + +import { getMasterDbClient, getReadReplicaDbClient } from '../database/client' import { createSettings } from './settings-factory' +import { DvmJobRepository } from '../repositories/dvm-job-repository' +import { DvmOrchestratorWorker } from '../app/dvm-orchestrator-worker' +import { EventRepository } from '../repositories/event-repository' export const dvmOrchestratorWorkerFactory = () => { - return new DvmOrchestratorWorker(process, createSettings) + const dbClient = getMasterDbClient() + const readReplicaDbClient = getReadReplicaDbClient() + const dvmJobRepository = new DvmJobRepository(dbClient) + const eventRepository = new EventRepository(dbClient, readReplicaDbClient, createSettings) + + return new DvmOrchestratorWorker(process, createSettings, dvmJobRepository, eventRepository) } diff --git a/src/repositories/dvm-job-repository.ts b/src/repositories/dvm-job-repository.ts index 89b0c15a..3a899c63 100644 --- a/src/repositories/dvm-job-repository.ts +++ b/src/repositories/dvm-job-repository.ts @@ -122,14 +122,23 @@ export class DvmJobRepository implements IDvmJobRepository { return row ? fromDBDvmJob(row) : undefined } - public async findPendingJobs(limit = 100, client: DatabaseClient = this.dbClient): Promise { - logger('find pending dvm jobs (limit %d)', limit) + public async findPendingJobs( + limit = 100, + kinds?: number[], + client: DatabaseClient = this.dbClient, + ): Promise { + logger('find pending dvm jobs (limit %d, kinds %o)', limit, kinds) - const rows = await client('dvm_jobs') + const query = client('dvm_jobs') .whereIn('status', [DvmJobStatus.SUBMITTED, DvmJobStatus.PICKED_UP]) .orderBy('created_at', 'asc') .limit(limit) - .select() + + if (Array.isArray(kinds) && kinds.length) { + query.whereIn('kind', kinds) + } + + const rows = await query.select() return rows.map(fromDBDvmJob) } diff --git a/test/unit/app/dvm-orchestrator-worker.spec.ts b/test/unit/app/dvm-orchestrator-worker.spec.ts index f53cc9cd..9e5bffd5 100644 --- a/test/unit/app/dvm-orchestrator-worker.spec.ts +++ b/test/unit/app/dvm-orchestrator-worker.spec.ts @@ -3,29 +3,104 @@ import EventEmitter from 'events' import Sinon from 'sinon' import sinonChai from 'sinon-chai' +import { DvmJobStatus } from '../../../src/@types/dvm' +import { IDvmJobRepository } from '../../../src/@types/repositories' import { Settings } from '../../../src/@types/settings' import { DvmOrchestratorWorker } from '../../../src/app/dvm-orchestrator-worker' +import * as processUtils from '../../../src/cli/utils/process' import * as metricsTelemetry from '../../../src/telemetry/metrics' +import * as eventUtils from '../../../src/utils/event' chai.use(sinonChai) const { expect } = chai +type FakeWorkerHandle = { + send: Sinon.SinonStub + kill: Sinon.SinonStub + onMessage: (handler: (message: unknown) => void) => void + onExit: (handler: (code: number | null, signal: NodeJS.Signals | null) => void) => void + onSpawnError: (handler: (reason: string) => void) => void + messageHandler?: (message: unknown) => void + exitHandler?: (code: number | null, signal: NodeJS.Signals | null) => void + spawnErrorHandler?: (reason: string) => void +} + +const flushMicrotasks = async (ticks = 10): Promise => { + for (let i = 0; i < ticks; i += 1) { + await Promise.resolve() + } +} + +const createFakeWorkerHandle = (): FakeWorkerHandle => { + const handle: FakeWorkerHandle = { + send: Sinon.stub().returns(true), + kill: Sinon.stub(), + onMessage: (handler) => { + handle.messageHandler = handler + }, + onExit: (handler) => { + handle.exitHandler = handler + }, + onSpawnError: (handler) => { + handle.spawnErrorHandler = handler + }, + } + return handle +} + describe('DvmOrchestratorWorker', () => { let sandbox: Sinon.SinonSandbox let fakeProcess: EventEmitter & { exit: Sinon.SinonStub; env: Record } let settings: Sinon.SinonStub let settingsState: Settings + let dvmJobRepository: Sinon.SinonStubbedInstance + let eventRepository: { findByFilters: Sinon.SinonStub; create: Sinon.SinonStub } + let spawnWorkerProcessStub: Sinon.SinonStub + let spawnedHandles: FakeWorkerHandle[] + let worker: DvmOrchestratorWorker | undefined + + const job = { + id: 'a'.repeat(64), + requesterPubkey: 'b'.repeat(64), + kind: 5000, + workerIndex: null, + status: DvmJobStatus.SUBMITTED, + resultEventId: null, + error: null, + pickedUpAt: null, + completedAt: null, + createdAt: new Date(), + updatedAt: new Date(), + } + + const requestEvent = { + id: job.id, + pubkey: job.requesterPubkey, + kind: 5000, + created_at: 1700000000, + tags: [], + content: 'do the job', + sig: 'sig', + } + + const currentWorkerHandle = () => spawnedHandles[spawnedHandles.length - 1] + + const createWorker = () => + new DvmOrchestratorWorker(fakeProcess as any, settings as any, dvmJobRepository as any, eventRepository as any) beforeEach(() => { sandbox = Sinon.createSandbox() + worker = undefined + spawnedHandles = [] fakeProcess = Object.assign(new EventEmitter(), { exit: sandbox.stub(), - env: {}, + env: { DVM_WORKER_INDEX: '0' }, }) as EventEmitter & { exit: Sinon.SinonStub; env: Record } settingsState = { + info: { relay_url: 'wss://relay.test' }, dvm: { workers: [{ command: 'python3', args: ['worker.py'] }], }, @@ -33,36 +108,257 @@ describe('DvmOrchestratorWorker', () => { settings = sandbox.stub().callsFake(() => settingsState) + dvmJobRepository = { + create: sandbox.stub(), + findById: sandbox.stub(), + assignWorker: sandbox.stub(), + updateStatus: sandbox.stub(), + findPendingJobs: sandbox.stub().resolves([]), + } as any + + eventRepository = { + findByFilters: sandbox.stub().resolves([{}]), + create: sandbox.stub().resolves(1), + } + + spawnWorkerProcessStub = sandbox.stub(processUtils, 'spawnWorkerProcess').callsFake(() => { + const handle = createFakeWorkerHandle() + spawnedHandles.push(handle) + return handle as any + }) + + sandbox.stub(eventUtils, 'toNostrEvent').returns(requestEvent as any) + sandbox.stub(eventUtils, 'getRelayPrivateKey').returns('relay-privkey') + sandbox.stub(eventUtils, 'getPublicKey').returns('relay-pubkey') + sandbox.stub(eventUtils, 'identifyEvent').callsFake(async (event: any) => ({ ...event, id: 'result-id' })) + sandbox.stub(eventUtils, 'signEvent').returns((event: any) => Promise.resolve({ ...event, sig: 'result-sig' })) + sandbox.stub(eventUtils, 'broadcastEvent').callsFake(async (event: any) => event) + sandbox.stub(metricsTelemetry, 'shutdownMetricsTelemetry').resolves() }) afterEach(() => { + worker?.close() 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) + it('logs startup and spawns the configured worker process', () => { + worker = createWorker() - expect(() => worker.run()).to.not.throw() + expect(() => worker!.run()).to.not.throw() expect(fakeProcess.exit).not.to.have.been.called + expect(spawnWorkerProcessStub).to.have.been.calledWith('python3', ['worker.py']) }) 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 = createWorker() worker.run() expect(fakeProcess.exit).to.have.been.calledWith(1) + expect(spawnWorkerProcessStub).not.to.have.been.called + }) + }) + + describe('dispatch', () => { + let clock: Sinon.SinonFakeTimers + + beforeEach(() => { + clock = Sinon.useFakeTimers() + }) + + afterEach(() => { + clock.restore() + }) + + it('does nothing when there are no pending jobs', async () => { + dvmJobRepository.findPendingJobs.resolves([]) + worker = createWorker() + worker.run() + + await clock.tickAsync(2000) + + expect(dvmJobRepository.assignWorker).not.to.have.been.called + expect(currentWorkerHandle().send).not.to.have.been.called + }) + + it('filters pending jobs by the worker config kinds', async () => { + settingsState.dvm = { workers: [{ command: 'python3', args: ['worker.py'], kinds: [5000, 5001] }] } + dvmJobRepository.findPendingJobs.resolves([]) + worker = createWorker() + worker.run() + + await clock.tickAsync(2000) + + expect(dvmJobRepository.findPendingJobs).to.have.been.calledWith(10, [5000, 5001]) + }) + + it('tries the next candidate in the batch when an earlier one cannot be assigned', async () => { + const job2 = { ...job, id: 'c'.repeat(64) } + dvmJobRepository.findPendingJobs.resolves([job, job2]) + dvmJobRepository.assignWorker.onFirstCall().resolves(false) + dvmJobRepository.assignWorker.onSecondCall().resolves(true) + worker = createWorker() + worker.run() + + await clock.tickAsync(2000) + + expect(dvmJobRepository.assignWorker).to.have.been.calledTwice + expect(dvmJobRepository.assignWorker.firstCall).to.have.been.calledWith(job.id, 0) + expect(dvmJobRepository.assignWorker.secondCall).to.have.been.calledWith(job2.id, 0) + expect(currentWorkerHandle().send).to.have.been.called + }) + + it('does not dispatch when it loses the race to assign the job', async () => { + dvmJobRepository.findPendingJobs.resolves([job]) + dvmJobRepository.assignWorker.resolves(false) + worker = createWorker() + worker.run() + + await clock.tickAsync(2000) + + expect(currentWorkerHandle().send).not.to.have.been.called + }) + + it('fails the job when the source event cannot be found', async () => { + dvmJobRepository.findPendingJobs.resolves([job]) + dvmJobRepository.assignWorker.resolves(true) + eventRepository.findByFilters.resolves([]) + worker = createWorker() + worker.run() + + await clock.tickAsync(2000) + + expect(dvmJobRepository.updateStatus).to.have.been.calledWithMatch({ + id: job.id, + status: DvmJobStatus.FAILED, + }) + expect(currentWorkerHandle().send).not.to.have.been.called + }) + + it('sends the job to the worker process as JSON over its stdin pipe', async () => { + dvmJobRepository.findPendingJobs.resolves([job]) + dvmJobRepository.assignWorker.resolves(true) + worker = createWorker() + worker.run() + + await clock.tickAsync(2000) + + expect(currentWorkerHandle().send).to.have.been.calledWith( + Sinon.match({ id: requestEvent.id, kind: requestEvent.kind }), + ) + }) + + it('fails the job instead of leaving it stuck when the worker pipe write fails', async () => { + dvmJobRepository.findPendingJobs.resolves([job]) + dvmJobRepository.assignWorker.resolves(true) + worker = createWorker() + worker.run() + + currentWorkerHandle().send.returns(false) + await clock.tickAsync(2000) + + expect(dvmJobRepository.updateStatus).to.have.been.calledWithMatch({ + id: job.id, + status: DvmJobStatus.FAILED, + }) + + expect(() => currentWorkerHandle().messageHandler!({ id: requestEvent.id, content: 'late reply' })).to.not.throw() + expect(eventRepository.create).not.to.have.been.called + }) + + it('marks the job completed and publishes a result event when the worker replies', async () => { + dvmJobRepository.findPendingJobs.resolves([job]) + dvmJobRepository.assignWorker.resolves(true) + worker = createWorker() + worker.run() + + await clock.tickAsync(2000) + currentWorkerHandle().messageHandler!({ id: requestEvent.id, content: 'result text' }) + await flushMicrotasks() + + expect(eventRepository.create).to.have.been.calledWithMatch({ kind: 6000, content: 'result text' }) + expect(dvmJobRepository.updateStatus).to.have.been.calledWithMatch({ + id: job.id, + status: DvmJobStatus.COMPLETED, + resultEventId: 'result-id', + }) + }) + + it('ignores worker replies for jobs that are no longer pending', async () => { + dvmJobRepository.findPendingJobs.resolves([]) + worker = createWorker() + worker.run() + + await clock.tickAsync(2000) + + expect(() => currentWorkerHandle().messageHandler!({ id: 'unknown-job-id', content: 'x' })).to.not.throw() + expect(dvmJobRepository.updateStatus).not.to.have.been.called + }) + + it('marks the job timed_out and kills the worker when no reply arrives before the configured timeout', async () => { + settingsState.dvm = { workers: [{ command: 'python3', args: ['worker.py'], timeoutMs: 5000 }] } + dvmJobRepository.findPendingJobs.resolves([job]) + dvmJobRepository.assignWorker.resolves(true) + worker = createWorker() + worker.run() + + await clock.tickAsync(2000) + const handle = currentWorkerHandle() + await clock.tickAsync(5000) + + expect(dvmJobRepository.updateStatus).to.have.been.calledWithMatch({ + id: job.id, + status: DvmJobStatus.TIMED_OUT, + }) + expect(handle.kill).to.have.been.calledOnce + }) + + it('marks all in-flight jobs failed when the worker process exits unexpectedly', async () => { + dvmJobRepository.findPendingJobs.resolves([job]) + dvmJobRepository.assignWorker.resolves(true) + worker = createWorker() + worker.run() + + await clock.tickAsync(2000) + currentWorkerHandle().exitHandler!(1, null) + + expect(dvmJobRepository.updateStatus).to.have.been.calledWithMatch({ + id: job.id, + status: DvmJobStatus.FAILED, + }) + }) + + it('respawns the worker process on the next tick after it exits', async () => { + dvmJobRepository.findPendingJobs.resolves([]) + worker = createWorker() + worker.run() + + expect(spawnedHandles).to.have.lengthOf(1) + spawnedHandles[0].exitHandler!(1, null) + + await clock.tickAsync(2000) + + expect(spawnedHandles).to.have.lengthOf(2) + }) + + it('does not dispatch a second job while a previous dispatch is still running', async () => { + dvmJobRepository.findPendingJobs.returns(new Promise(() => {})) // never resolves + worker = createWorker() + worker.run() + + await clock.tickAsync(2000) + await clock.tickAsync(2000) + + expect(dvmJobRepository.findPendingJobs).to.have.been.calledOnce }) }) 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 = createWorker() worker.run() fakeProcess.emit('SIGTERM') @@ -75,13 +371,40 @@ describe('DvmOrchestratorWorker', () => { }) describe('close', () => { - it('invokes the callback', () => { - const worker = new DvmOrchestratorWorker(fakeProcess as any, settings as any) + it('invokes the callback and kills the worker process', () => { + worker = createWorker() + worker.run() const callback = sandbox.stub() worker.close(callback) expect(callback).to.have.been.calledOnce + expect(currentWorkerHandle().kill).to.have.been.calledOnce + }) + + it('fails in-flight jobs before invoking the callback', async () => { + const clock = Sinon.useFakeTimers() + try { + dvmJobRepository.findPendingJobs.resolves([job]) + dvmJobRepository.assignWorker.resolves(true) + worker = createWorker() + worker.run() + + await clock.tickAsync(2000) + + const callback = sandbox.stub() + worker.close(callback) + await flushMicrotasks() + + expect(dvmJobRepository.updateStatus).to.have.been.calledWithMatch({ + id: job.id, + status: DvmJobStatus.FAILED, + }) + expect(callback).to.have.been.calledOnce + expect(currentWorkerHandle().kill).to.have.been.calledOnce + } finally { + clock.restore() + } }) }) }) diff --git a/test/unit/cli/process.spec.ts b/test/unit/cli/process.spec.ts new file mode 100644 index 00000000..c5e6518d --- /dev/null +++ b/test/unit/cli/process.spec.ts @@ -0,0 +1,151 @@ +import chai from 'chai' +import EventEmitter from 'events' +import Sinon from 'sinon' +import sinonChai from 'sinon-chai' + +import { spawnWorkerProcess } from '../../../src/cli/utils/process' + +// `import * as childProcess` goes through TS's __importStar interop helper, +// which wraps built-in CJS modules in a frozen getter-only object sinon can't +// stub. `import ... = require(...)` gives the raw module object instead — +// the same one process.ts's own `require('child_process')` call resolves to. +import childProcess = require('child_process') + +chai.use(sinonChai) + +const { expect } = chai + +type FakeChildProcess = EventEmitter & { + stdout: EventEmitter + stderr: EventEmitter + stdin: { write: Sinon.SinonStub } + kill: Sinon.SinonStub +} + +const createFakeChild = (): FakeChildProcess => + Object.assign(new EventEmitter(), { + stdout: new EventEmitter(), + stderr: new EventEmitter(), + stdin: { write: Sinon.stub() }, + kill: Sinon.stub(), + }) + +describe('spawnWorkerProcess', () => { + let sandbox: Sinon.SinonSandbox + let spawnStub: Sinon.SinonStub + let fakeChild: FakeChildProcess + + beforeEach(() => { + sandbox = Sinon.createSandbox() + fakeChild = createFakeChild() + spawnStub = sandbox.stub(childProcess, 'spawn').returns(fakeChild as any) + }) + + afterEach(() => { + sandbox.restore() + }) + + it('spawns the command with piped stdio and no shell', () => { + spawnWorkerProcess('python3', ['worker.py']) + + expect(spawnStub).to.have.been.calledWith('python3', ['worker.py'], Sinon.match({ stdio: 'pipe', shell: false })) + }) + + it('writes sent messages to stdin as a newline-terminated JSON line', () => { + const handle = spawnWorkerProcess('python3', ['worker.py']) + + const sent = handle.send({ id: 'abc', kind: 5000 }) + + expect(fakeChild.stdin.write).to.have.been.calledWith('{"id":"abc","kind":5000}\n') + expect(sent).to.be.true + }) + + it('returns false instead of throwing when stdin.write fails on a dead worker', () => { + fakeChild.stdin.write.throws(new Error('write after end')) + const handle = spawnWorkerProcess('python3', ['worker.py']) + + const sent = handle.send({ id: 'abc', kind: 5000 }) + + expect(sent).to.be.false + }) + + it('parses newline-delimited JSON from stdout and dispatches each line to message handlers', () => { + const handle = spawnWorkerProcess('python3', ['worker.py']) + const received: unknown[] = [] + handle.onMessage((message) => received.push(message)) + + fakeChild.stdout.emit('data', Buffer.from('{"id":"a"}\n{"id":"b"}\n')) + + expect(received).to.deep.equal([{ id: 'a' }, { id: 'b' }]) + }) + + it('buffers a line split across multiple stdout chunks', () => { + const handle = spawnWorkerProcess('python3', ['worker.py']) + const received: unknown[] = [] + handle.onMessage((message) => received.push(message)) + + fakeChild.stdout.emit('data', Buffer.from('{"id":"a"')) + fakeChild.stdout.emit('data', Buffer.from('}\n')) + + expect(received).to.deep.equal([{ id: 'a' }]) + }) + + it('drops malformed JSON lines without throwing', () => { + const handle = spawnWorkerProcess('python3', ['worker.py']) + const received: unknown[] = [] + handle.onMessage((message) => received.push(message)) + + expect(() => fakeChild.stdout.emit('data', Buffer.from('not json\n{"id":"a"}\n'))).to.not.throw() + expect(received).to.deep.equal([{ id: 'a' }]) + }) + + it('classifies ENOENT spawn errors as not-found', () => { + const handle = spawnWorkerProcess('nope', []) + const reasons: string[] = [] + handle.onSpawnError((reason) => reasons.push(reason)) + + const error = Object.assign(new Error('not found'), { code: 'ENOENT' }) + fakeChild.emit('error', error) + + expect(reasons).to.deep.equal(['not-found']) + }) + + it('classifies EACCES spawn errors as permission-denied', () => { + const handle = spawnWorkerProcess('nope', []) + const reasons: string[] = [] + handle.onSpawnError((reason) => reasons.push(reason)) + + const error = Object.assign(new Error('denied'), { code: 'EACCES' }) + fakeChild.emit('error', error) + + expect(reasons).to.deep.equal(['permission-denied']) + }) + + it('classifies other spawn errors as spawn-error', () => { + const handle = spawnWorkerProcess('nope', []) + const reasons: string[] = [] + handle.onSpawnError((reason) => reasons.push(reason)) + + fakeChild.emit('error', new Error('boom')) + + expect(reasons).to.deep.equal(['spawn-error']) + }) + + it('notifies exit handlers with the exit code and signal', () => { + const handle = spawnWorkerProcess('python3', ['worker.py']) + const exits: Array<[number | null, NodeJS.Signals | null]> = [] + handle.onExit((code, signal) => exits.push([code, signal])) + + fakeChild.emit('exit', 1, null) + + expect(exits).to.deep.equal([[1, null]]) + }) + + it('kills the underlying child process with SIGTERM', () => { + const handle = spawnWorkerProcess('python3', ['worker.py']) + + handle.kill() + + expect(fakeChild.kill).to.have.been.calledWith('SIGTERM') + }) +}) diff --git a/test/unit/repositories/dvm-job-repository.spec.ts b/test/unit/repositories/dvm-job-repository.spec.ts index c23236b6..c8db9942 100644 --- a/test/unit/repositories/dvm-job-repository.spec.ts +++ b/test/unit/repositories/dvm-job-repository.spec.ts @@ -295,7 +295,7 @@ describe('DvmJobRepository', () => { const whereInStub = sandbox.stub().returns({ orderBy: orderByStub }) const client = sandbox.stub().returns({ whereIn: whereInStub }) as unknown as DatabaseClient - const result = await repository.findPendingJobs(10, client) + const result = await repository.findPendingJobs(10, undefined, client) expect(result).to.be.an('array').that.is.empty }) @@ -307,7 +307,7 @@ describe('DvmJobRepository', () => { const whereInStub = sandbox.stub().returns({ orderBy: orderByStub }) const client = sandbox.stub().returns({ whereIn: whereInStub }) as unknown as DatabaseClient - const result = await repository.findPendingJobs(10, client) + const result = await repository.findPendingJobs(10, undefined, client) expect(result).to.have.lengthOf(1) expect(result[0].id).to.equal(jobId) @@ -320,7 +320,7 @@ describe('DvmJobRepository', () => { const whereInStub = sandbox.stub().returns({ orderBy: orderByStub }) const client = sandbox.stub().returns({ whereIn: whereInStub }) as unknown as DatabaseClient - await repository.findPendingJobs(10, client) + await repository.findPendingJobs(10, undefined, client) expect(whereInStub).to.have.been.calledWith('status', [DvmJobStatus.SUBMITTED, DvmJobStatus.PICKED_UP]) }) @@ -332,7 +332,7 @@ describe('DvmJobRepository', () => { const whereInStub = sandbox.stub().returns({ orderBy: orderByStub }) const client = sandbox.stub().returns({ whereIn: whereInStub }) as unknown as DatabaseClient - await repository.findPendingJobs(10, client) + await repository.findPendingJobs(10, undefined, client) expect(orderByStub).to.have.been.calledWith('created_at', 'asc') }) @@ -344,9 +344,34 @@ describe('DvmJobRepository', () => { const whereInStub = sandbox.stub().returns({ orderBy: orderByStub }) const client = sandbox.stub().returns({ whereIn: whereInStub }) as unknown as DatabaseClient - await repository.findPendingJobs(undefined, client) + await repository.findPendingJobs(undefined, undefined, client) expect(limitStub).to.have.been.calledWith(100) }) + + it('filters by kind when kinds are provided', async () => { + const selectStub = sandbox.stub().resolves([]) + const kindWhereInStub = sandbox.stub() + const limitStub = sandbox.stub().returns({ select: selectStub, whereIn: kindWhereInStub }) + const orderByStub = sandbox.stub().returns({ limit: limitStub }) + const whereInStub = sandbox.stub().returns({ orderBy: orderByStub }) + const client = sandbox.stub().returns({ whereIn: whereInStub }) as unknown as DatabaseClient + + await repository.findPendingJobs(10, [5000, 5001], client) + + expect(kindWhereInStub).to.have.been.calledWith('kind', [5000, 5001]) + }) + + it('does not filter by kind when kinds is omitted or empty', async () => { + const selectStub = sandbox.stub().resolves([]) + const limitStub = sandbox.stub().returns({ select: selectStub }) + const orderByStub = sandbox.stub().returns({ limit: limitStub }) + const whereInStub = sandbox.stub().returns({ orderBy: orderByStub }) + const client = sandbox.stub().returns({ whereIn: whereInStub }) as unknown as DatabaseClient + + const result = await repository.findPendingJobs(10, [], client) + + expect(result).to.be.an('array').that.is.empty + }) }) })