From 88860e4f8e0c5d3bb370821088b2e873aee23bf2 Mon Sep 17 00:00:00 2001 From: Priyanshubhartistm Date: Wed, 12 Aug 2026 16:03:15 +0530 Subject: [PATCH 1/3] feat(dvm): add dvm job persistence migration and repository Signed-off-by: Priyanshubhartistm --- .changeset/dvm-job-persistence.md | 5 + .knip.json | 1 + .../20260812_150000_create_dvm_jobs_table.js | 23 +++ src/@types/dvm.ts | 37 +++++ src/@types/repositories.ts | 16 ++- src/repositories/dvm-job-repository.ts | 133 ++++++++++++++++++ 6 files changed, 212 insertions(+), 3 deletions(-) create mode 100644 .changeset/dvm-job-persistence.md create mode 100644 migrations/20260812_150000_create_dvm_jobs_table.js create mode 100644 src/@types/dvm.ts create mode 100644 src/repositories/dvm-job-repository.ts diff --git a/.changeset/dvm-job-persistence.md b/.changeset/dvm-job-persistence.md new file mode 100644 index 00000000..c71f4f0a --- /dev/null +++ b/.changeset/dvm-job-persistence.md @@ -0,0 +1,5 @@ +--- +"nostream": minor +--- + +feat(dvm): add job persistence migration and repository for DVM job state diff --git a/.knip.json b/.knip.json index 7129b278..f47eb4b9 100644 --- a/.knip.json +++ b/.knip.json @@ -17,6 +17,7 @@ "ignore": [ ".nostr/**", "src/repositories/invite-code-repository.ts", + "src/repositories/dvm-job-repository.ts", "src/utils/relay-probe/**" ], "commitlint": false, diff --git a/migrations/20260812_150000_create_dvm_jobs_table.js b/migrations/20260812_150000_create_dvm_jobs_table.js new file mode 100644 index 00000000..6bbb1e56 --- /dev/null +++ b/migrations/20260812_150000_create_dvm_jobs_table.js @@ -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') +} diff --git a/src/@types/dvm.ts b/src/@types/dvm.ts new file mode 100644 index 00000000..0f5c5940 --- /dev/null +++ b/src/@types/dvm.ts @@ -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 +} diff --git a/src/@types/repositories.ts b/src/@types/repositories.ts index df98e545..717ba034 100644 --- a/src/@types/repositories.ts +++ b/src/@types/repositories.ts @@ -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' @@ -73,3 +73,13 @@ export interface IInviteCodeRepository { findActiveCodes(limit?: number): Promise deleteExpiredCodes(): Promise } + +export interface IDvmJobRepository { + create(id: string, requesterPubkey: Pubkey, kind: number): Promise + findById(id: string): Promise + assignWorker(id: string, workerIndex: number): Promise + updateStatus( + job: Pick & Partial>, + ): Promise + findPendingJobs(limit?: number): Promise +} diff --git a/src/repositories/dvm-job-repository.ts b/src/repositories/dvm-job-repository.ts new file mode 100644 index 00000000..32dbe629 --- /dev/null +++ b/src/repositories/dvm-job-repository.ts @@ -0,0 +1,133 @@ +import { DatabaseClient, Pubkey } from '../@types/base' +import { DBDvmJob, DvmJob, DvmJobStatus } from '../@types/dvm' +import { IDvmJobRepository } from '../@types/repositories' +import { createLogger } from '../factories/logger-factory' +import { fromBuffer, toBuffer } from '../utils/transform' + +const logger = createLogger('dvm-job-repository') + +function fromDBDvmJob(row: DBDvmJob): DvmJob { + return { + id: fromBuffer(row.id), + requesterPubkey: fromBuffer(row.requester_pubkey), + kind: row.kind, + workerIndex: row.worker_index, + status: row.status, + resultEventId: row.result_event_id ? fromBuffer(row.result_event_id) : null, + error: row.error, + pickedUpAt: row.picked_up_at, + completedAt: row.completed_at, + createdAt: row.created_at, + updatedAt: row.updated_at, + } +} + +function affectedRows(result: unknown): number { + if (typeof result === 'number') { + return result + } + if (result && typeof (result as any).rowCount === 'number') { + return (result as any).rowCount + } + return 0 +} + +export class DvmJobRepository implements IDvmJobRepository { + public constructor(private readonly dbClient: DatabaseClient) {} + + public async create( + id: string, + requesterPubkey: Pubkey, + kind: number, + client: DatabaseClient = this.dbClient, + ): Promise { + logger('create dvm job %s (kind %d) for %s', id, kind, requesterPubkey) + + const now = new Date() + const row: DBDvmJob = { + id: toBuffer(id), + requester_pubkey: toBuffer(requesterPubkey), + kind, + worker_index: null, + status: DvmJobStatus.SUBMITTED, + result_event_id: null, + error: null, + picked_up_at: null, + completed_at: null, + created_at: now, + updated_at: now, + } + + await client('dvm_jobs').insert(row) + + return fromDBDvmJob(row) + } + + public async findById(id: string, client: DatabaseClient = this.dbClient): Promise { + logger('find dvm job %s', id) + + const [row] = await client('dvm_jobs').where('id', toBuffer(id)).select() + + if (!row) { + return + } + + return fromDBDvmJob(row) + } + + // Atomic pickup: single conditional UPDATE ensures only one worker wins the job + public async assignWorker(id: string, workerIndex: number, client: DatabaseClient = this.dbClient): Promise { + logger('assign dvm job %s to worker %d', id, workerIndex) + + const now = new Date() + + const result = await client('dvm_jobs') + .where('id', toBuffer(id)) + .where('status', DvmJobStatus.SUBMITTED) + .update({ + worker_index: workerIndex, + status: DvmJobStatus.PICKED_UP, + picked_up_at: now, + updated_at: now, + }) + + return affectedRows(result) > 0 + } + + public async updateStatus( + job: Pick & Partial>, + client: DatabaseClient = this.dbClient, + ): Promise { + logger('update dvm job status: %o', job) + + const now = new Date() + const isTerminal = + job.status === DvmJobStatus.COMPLETED || + job.status === DvmJobStatus.FAILED || + job.status === DvmJobStatus.TIMED_OUT + + const update: Partial = { + status: job.status, + updated_at: now, + ...(isTerminal ? { completed_at: now } : {}), + ...(job.resultEventId ? { result_event_id: toBuffer(job.resultEventId) } : {}), + ...(job.error ? { error: job.error } : {}), + } + + const [row] = await client('dvm_jobs').where('id', toBuffer(job.id)).update(update).returning(['*']) + + return row ? fromDBDvmJob(row) : undefined + } + + public async findPendingJobs(limit = 100, client: DatabaseClient = this.dbClient): Promise { + logger('find pending dvm jobs (limit %d)', limit) + + const rows = await client('dvm_jobs') + .whereIn('status', [DvmJobStatus.SUBMITTED, DvmJobStatus.PICKED_UP]) + .orderBy('created_at', 'asc') + .limit(limit) + .select() + + return rows.map(fromDBDvmJob) + } +} From 7de3528de682644421fa99f3371c1c93f65ae486 Mon Sep 17 00:00:00 2001 From: Priyanshubhartistm Date: Wed, 12 Aug 2026 16:04:11 +0530 Subject: [PATCH 2/3] test(dvm): add dvm-job-repository unit tests Signed-off-by: Priyanshubhartistm --- .../repositories/dvm-job-repository.spec.ts | 315 ++++++++++++++++++ 1 file changed, 315 insertions(+) create mode 100644 test/unit/repositories/dvm-job-repository.spec.ts diff --git a/test/unit/repositories/dvm-job-repository.spec.ts b/test/unit/repositories/dvm-job-repository.spec.ts new file mode 100644 index 00000000..3bd1ef59 --- /dev/null +++ b/test/unit/repositories/dvm-job-repository.spec.ts @@ -0,0 +1,315 @@ +import * as chai from 'chai' +import chaiAsPromised from 'chai-as-promised' +import * as sinon from 'sinon' +import sinonChai from 'sinon-chai' + +import { DatabaseClient } from '../../../src/@types/base' +import { DvmJobStatus } from '../../../src/@types/dvm' +import { DvmJobRepository } from '../../../src/repositories/dvm-job-repository' + +chai.use(sinonChai) +chai.use(chaiAsPromised) + +const { expect } = chai + +describe('DvmJobRepository', () => { + let repository: DvmJobRepository + let sandbox: sinon.SinonSandbox + + const fixedDate = new Date('2026-08-12T00:00:00.000Z') + const jobId = 'a'.repeat(64) + const pubkeyHex = '22e804d26ed16b68db5259e78449e96dab5d464c8f470bda3eb1a70467f2c793' + const resultEventId = 'b'.repeat(64) + + const dbDvmJobRow = { + id: Buffer.from(jobId, 'hex'), + requester_pubkey: Buffer.from(pubkeyHex, 'hex'), + kind: 5000, + worker_index: null as number | null, + status: DvmJobStatus.SUBMITTED, + result_event_id: null as Buffer | null, + error: null as string | null, + picked_up_at: null as Date | null, + completed_at: null as Date | null, + created_at: fixedDate, + updated_at: fixedDate, + } + + beforeEach(() => { + sandbox = sinon.createSandbox() + sandbox.useFakeTimers(fixedDate.getTime()) + + repository = new DvmJobRepository({} as DatabaseClient) + }) + + afterEach(() => { + sandbox.restore() + }) + + describe('.create', () => { + it('inserts into the dvm_jobs table', async () => { + const insertStub = sandbox.stub().resolves() + const client = sandbox.stub().returns({ + insert: insertStub, + }) as unknown as DatabaseClient + + await repository.create(jobId, pubkeyHex, 5000, client) + + expect(client).to.have.been.calledWith('dvm_jobs') + }) + + it('returns a DvmJob with submitted status and no worker assigned', async () => { + const insertStub = sandbox.stub().resolves() + const client = sandbox.stub().returns({ + insert: insertStub, + }) as unknown as DatabaseClient + + const result = await repository.create(jobId, pubkeyHex, 5000, client) + + expect(result).to.deep.include({ + id: jobId, + requesterPubkey: pubkeyHex, + kind: 5000, + workerIndex: null, + status: DvmJobStatus.SUBMITTED, + resultEventId: null, + error: null, + }) + expect(result.createdAt).to.be.instanceOf(Date) + expect(result.updatedAt).to.be.instanceOf(Date) + }) + + it('stores id and requester pubkey as buffers', async () => { + const insertStub = sandbox.stub().resolves() + const client = sandbox.stub().returns({ + insert: insertStub, + }) as unknown as DatabaseClient + + await repository.create(jobId, pubkeyHex, 5000, client) + + const insertedRow = insertStub.firstCall.args[0] + expect(insertedRow.id).to.deep.equal(Buffer.from(jobId, 'hex')) + expect(insertedRow.requester_pubkey).to.deep.equal(Buffer.from(pubkeyHex, 'hex')) + }) + }) + + describe('.findById', () => { + it('returns undefined when no job is found', async () => { + const client = sandbox.stub().returns({ + where: sandbox.stub().returns({ select: sandbox.stub().resolves([]) }), + }) as unknown as DatabaseClient + + const result = await repository.findById(jobId, client) + + expect(result).to.be.undefined + }) + + it('returns a transformed DvmJob when found', async () => { + const client = sandbox.stub().returns({ + where: sandbox.stub().returns({ select: sandbox.stub().resolves([dbDvmJobRow]) }), + }) as unknown as DatabaseClient + + const result = await repository.findById(jobId, client) + + expect(result).to.not.be.undefined + expect(result!.id).to.equal(jobId) + expect(result!.requesterPubkey).to.equal(pubkeyHex) + expect(result!.status).to.equal(DvmJobStatus.SUBMITTED) + }) + + it('queries the dvm_jobs table by id', async () => { + const whereStub = sandbox.stub().returns({ select: sandbox.stub().resolves([]) }) + const client = sandbox.stub().returns({ where: whereStub }) as unknown as DatabaseClient + + await repository.findById(jobId, client) + + expect(client).to.have.been.calledWith('dvm_jobs') + const [field, value] = whereStub.firstCall.args + expect(field).to.equal('id') + expect(value).to.deep.equal(Buffer.from(jobId, 'hex')) + }) + }) + + describe('.assignWorker', () => { + it('returns true when assignment succeeds (rowCount > 0)', async () => { + const updateStub = sandbox.stub().resolves(1) + const whereStub2 = sandbox.stub().returns({ update: updateStub }) + const whereStub1 = sandbox.stub().returns({ where: whereStub2 }) + const client = sandbox.stub().returns({ where: whereStub1 }) as unknown as DatabaseClient + + const result = await repository.assignWorker(jobId, 0, client) + + expect(result).to.be.true + }) + + it('returns false when no submitted job matched (rowCount = 0)', async () => { + const updateStub = sandbox.stub().resolves(0) + const whereStub2 = sandbox.stub().returns({ update: updateStub }) + const whereStub1 = sandbox.stub().returns({ where: whereStub2 }) + const client = sandbox.stub().returns({ where: whereStub1 }) as unknown as DatabaseClient + + const result = await repository.assignWorker(jobId, 0, client) + + expect(result).to.be.false + }) + + it('returns true when pg returns { rowCount } object', async () => { + const updateStub = sandbox.stub().resolves({ rowCount: 1 }) + const whereStub2 = sandbox.stub().returns({ update: updateStub }) + const whereStub1 = sandbox.stub().returns({ where: whereStub2 }) + const client = sandbox.stub().returns({ where: whereStub1 }) as unknown as DatabaseClient + + const result = await repository.assignWorker(jobId, 0, client) + + expect(result).to.be.true + }) + + it('only matches jobs still in submitted status', async () => { + const updateStub = sandbox.stub().resolves(1) + const whereStub2 = sandbox.stub().returns({ update: updateStub }) + const whereStub1 = sandbox.stub().returns({ where: whereStub2 }) + const client = sandbox.stub().returns({ where: whereStub1 }) as unknown as DatabaseClient + + await repository.assignWorker(jobId, 2, client) + + expect(whereStub2).to.have.been.calledWith('status', DvmJobStatus.SUBMITTED) + }) + }) + + describe('.updateStatus', () => { + it('updates status and returns the transformed job', async () => { + const updatedRow = { ...dbDvmJobRow, status: DvmJobStatus.COMPLETED, completed_at: fixedDate } + const returningStub = sandbox.stub().resolves([updatedRow]) + const updateStub = sandbox.stub().returns({ returning: returningStub }) + const whereStub = sandbox.stub().returns({ update: updateStub }) + const client = sandbox.stub().returns({ where: whereStub }) as unknown as DatabaseClient + + const result = await repository.updateStatus({ id: jobId, status: DvmJobStatus.COMPLETED }, client) + + expect(result).to.not.be.undefined + expect(result!.status).to.equal(DvmJobStatus.COMPLETED) + }) + + it('returns undefined when no matching job exists', async () => { + const returningStub = sandbox.stub().resolves([]) + const updateStub = sandbox.stub().returns({ returning: returningStub }) + const whereStub = sandbox.stub().returns({ update: updateStub }) + const client = sandbox.stub().returns({ where: whereStub }) as unknown as DatabaseClient + + const result = await repository.updateStatus({ id: jobId, status: DvmJobStatus.FAILED }, client) + + expect(result).to.be.undefined + }) + + it('sets completed_at for terminal statuses', async () => { + const returningStub = sandbox.stub().resolves([dbDvmJobRow]) + const updateStub = sandbox.stub().returns({ returning: returningStub }) + const whereStub = sandbox.stub().returns({ update: updateStub }) + const client = sandbox.stub().returns({ where: whereStub }) as unknown as DatabaseClient + + await repository.updateStatus({ id: jobId, status: DvmJobStatus.TIMED_OUT }, client) + + const update = updateStub.firstCall.args[0] + expect(update.completed_at).to.deep.equal(fixedDate) + }) + + it('does not set completed_at for the picked_up status', async () => { + const returningStub = sandbox.stub().resolves([dbDvmJobRow]) + const updateStub = sandbox.stub().returns({ returning: returningStub }) + const whereStub = sandbox.stub().returns({ update: updateStub }) + const client = sandbox.stub().returns({ where: whereStub }) as unknown as DatabaseClient + + await repository.updateStatus({ id: jobId, status: DvmJobStatus.PICKED_UP }, client) + + const update = updateStub.firstCall.args[0] + expect(update.completed_at).to.be.undefined + }) + + it('encodes resultEventId as a buffer when provided', async () => { + const returningStub = sandbox.stub().resolves([dbDvmJobRow]) + const updateStub = sandbox.stub().returns({ returning: returningStub }) + const whereStub = sandbox.stub().returns({ update: updateStub }) + const client = sandbox.stub().returns({ where: whereStub }) as unknown as DatabaseClient + + await repository.updateStatus({ id: jobId, status: DvmJobStatus.COMPLETED, resultEventId }, client) + + const update = updateStub.firstCall.args[0] + expect(update.result_event_id).to.deep.equal(Buffer.from(resultEventId, 'hex')) + }) + + it('sets error when provided', async () => { + const returningStub = sandbox.stub().resolves([dbDvmJobRow]) + const updateStub = sandbox.stub().returns({ returning: returningStub }) + const whereStub = sandbox.stub().returns({ update: updateStub }) + const client = sandbox.stub().returns({ where: whereStub }) as unknown as DatabaseClient + + await repository.updateStatus({ id: jobId, status: DvmJobStatus.FAILED, error: 'worker crashed' }, client) + + const update = updateStub.firstCall.args[0] + expect(update.error).to.equal('worker crashed') + }) + }) + + describe('.findPendingJobs', () => { + it('returns an empty array when no pending jobs exist', 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 + }) + + it('returns transformed DvmJob objects', async () => { + const selectStub = sandbox.stub().resolves([dbDvmJobRow]) + 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.have.lengthOf(1) + expect(result[0].id).to.equal(jobId) + }) + + it('filters for submitted and picked_up statuses', 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 + + await repository.findPendingJobs(10, client) + + expect(whereInStub).to.have.been.calledWith('status', [DvmJobStatus.SUBMITTED, DvmJobStatus.PICKED_UP]) + }) + + it('orders by created_at ascending', 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 + + await repository.findPendingJobs(10, client) + + expect(orderByStub).to.have.been.calledWith('created_at', 'asc') + }) + + it('defaults limit to 100', 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 + + await repository.findPendingJobs(undefined, client) + + expect(limitStub).to.have.been.calledWith(100) + }) + }) +}) From 5c212bfee4eb4fb43bba06147b85c2dc0502423b Mon Sep 17 00:00:00 2001 From: Priyanshubhartistm Date: Wed, 12 Aug 2026 22:11:18 +0530 Subject: [PATCH 3/3] fix(dvm): use composite status/created_at index and fix null-clearing in updateStatus Signed-off-by: Priyanshubhartistm --- .../20260812_150000_create_dvm_jobs_table.js | 10 ++++- src/repositories/dvm-job-repository.ts | 7 +++- .../repositories/dvm-job-repository.spec.ts | 37 +++++++++++++++++++ 3 files changed, 50 insertions(+), 4 deletions(-) diff --git a/migrations/20260812_150000_create_dvm_jobs_table.js b/migrations/20260812_150000_create_dvm_jobs_table.js index 6bbb1e56..0e06e081 100644 --- a/migrations/20260812_150000_create_dvm_jobs_table.js +++ b/migrations/20260812_150000_create_dvm_jobs_table.js @@ -4,7 +4,10 @@ exports.up = function (knex) { 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 + .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() @@ -13,7 +16,10 @@ exports.up = function (knex) { 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') + // Composite (not status-only): findPendingJobs() filters by status AND + // orders by created_at, so the index needs to satisfy both the filter + // and the sort for FIFO polling, same as invoices_pending_created_at_idx. + table.index(['status', 'created_at'], 'idx_dvm_jobs_status_created_at') table.index(['kind'], 'idx_dvm_jobs_kind') }) } diff --git a/src/repositories/dvm-job-repository.ts b/src/repositories/dvm-job-repository.ts index 32dbe629..89b0c15a 100644 --- a/src/repositories/dvm-job-repository.ts +++ b/src/repositories/dvm-job-repository.ts @@ -106,12 +106,15 @@ export class DvmJobRepository implements IDvmJobRepository { job.status === DvmJobStatus.FAILED || job.status === DvmJobStatus.TIMED_OUT + // Check key presence, not truthiness: a caller passing `resultEventId: null` + // or `error: null` is explicitly clearing the field, which a truthy check + // would silently ignore and leave the stale DB value in place. const update: Partial = { status: job.status, updated_at: now, ...(isTerminal ? { completed_at: now } : {}), - ...(job.resultEventId ? { result_event_id: toBuffer(job.resultEventId) } : {}), - ...(job.error ? { error: job.error } : {}), + ...('resultEventId' in job ? { result_event_id: job.resultEventId ? toBuffer(job.resultEventId) : null } : {}), + ...('error' in job ? { error: job.error ?? null } : {}), } const [row] = await client('dvm_jobs').where('id', toBuffer(job.id)).update(update).returning(['*']) diff --git a/test/unit/repositories/dvm-job-repository.spec.ts b/test/unit/repositories/dvm-job-repository.spec.ts index 3bd1ef59..c23236b6 100644 --- a/test/unit/repositories/dvm-job-repository.spec.ts +++ b/test/unit/repositories/dvm-job-repository.spec.ts @@ -248,6 +248,43 @@ describe('DvmJobRepository', () => { const update = updateStub.firstCall.args[0] expect(update.error).to.equal('worker crashed') }) + + it('does not touch resultEventId or error when the keys are omitted', async () => { + const returningStub = sandbox.stub().resolves([dbDvmJobRow]) + const updateStub = sandbox.stub().returns({ returning: returningStub }) + const whereStub = sandbox.stub().returns({ update: updateStub }) + const client = sandbox.stub().returns({ where: whereStub }) as unknown as DatabaseClient + + await repository.updateStatus({ id: jobId, status: DvmJobStatus.PICKED_UP }, client) + + const update = updateStub.firstCall.args[0] + expect(update).to.not.have.property('result_event_id') + expect(update).to.not.have.property('error') + }) + + it('clears resultEventId when explicitly set to null', async () => { + const returningStub = sandbox.stub().resolves([dbDvmJobRow]) + const updateStub = sandbox.stub().returns({ returning: returningStub }) + const whereStub = sandbox.stub().returns({ update: updateStub }) + const client = sandbox.stub().returns({ where: whereStub }) as unknown as DatabaseClient + + await repository.updateStatus({ id: jobId, status: DvmJobStatus.PICKED_UP, resultEventId: null }, client) + + const update = updateStub.firstCall.args[0] + expect(update.result_event_id).to.be.null + }) + + it('clears error when explicitly set to null', async () => { + const returningStub = sandbox.stub().resolves([dbDvmJobRow]) + const updateStub = sandbox.stub().returns({ returning: returningStub }) + const whereStub = sandbox.stub().returns({ update: updateStub }) + const client = sandbox.stub().returns({ where: whereStub }) as unknown as DatabaseClient + + await repository.updateStatus({ id: jobId, status: DvmJobStatus.PICKED_UP, error: null }, client) + + const update = updateStub.firstCall.args[0] + expect(update.error).to.be.null + }) }) describe('.findPendingJobs', () => {