diff --git a/src/backend/common/database/drizzle/drizzleUtils.ts b/src/backend/common/database/drizzle/drizzleUtils.ts index 76814d7a..5fa2c74b 100644 --- a/src/backend/common/database/drizzle/drizzleUtils.ts +++ b/src/backend/common/database/drizzle/drizzleUtils.ts @@ -79,7 +79,7 @@ export const getDb = (dbVal: string, opts: { logger?: Logger, backupPath?: strin logger = loggerNoop, } = opts; const db = drizzle({relations: relations, logger: createDrizzleLogger(logger), connection: {path: dbVal, allowExtension: true}}); - if(dbVal !== ':memory:') { + if(dbVal !== MEMORY_DB_NAME) { logger.debug('Loading honker extension'); db.$client.loadExtension(extensionPath()); db.run('SELECT honker_bootstrap()'); diff --git a/src/backend/common/database/honker/HonkerQueue.ts b/src/backend/common/database/honker/HonkerQueue.ts new file mode 100644 index 00000000..f4ac3e85 --- /dev/null +++ b/src/backend/common/database/honker/HonkerQueue.ts @@ -0,0 +1,49 @@ +import { sql } from 'drizzle-orm'; +import type { DbConcrete } from '../drizzle/drizzleUtils.ts'; +import type {Job} from '@russellthehippo/honker-node'; + +export interface HonkerJobData { + id: number, + queue: string + payload: string + worker_id: string + attempts: number + claim_expires_at: number +} + +export class Queue { + private readonly name: string; + private readonly maxAttempts: number; + private readonly visibilityTimeout: number; + private db: DbConcrete; + constructor( + db: DbConcrete, + name: string, + maxAttempts: number = 3 + ) { + this.db = db; + this.name = name; + this.maxAttempts = maxAttempts; + } + + enqueue( + payload: T, + opts: { delay?: number; priority?: number, tx?: DbConcrete } = {}, + ): number { + const db = opts.tx ?? this.db; + const row = db.get<{ id: number }>(sql` + SELECT honker_enqueue( + ${this.name}, + ${JSON.stringify(payload)}, + NULL, + ${opts.delay ?? null}, + ${opts.priority ?? 0}, + ${this.maxAttempts}, + NULL + ) AS id + `); + return row!.id; + } +} + +export type HonkerJob = Omit & {payload: T}; \ No newline at end of file diff --git a/src/backend/tests/honker/honker.test.ts b/src/backend/tests/honker/honker.test.ts new file mode 100644 index 00000000..39f29755 --- /dev/null +++ b/src/backend/tests/honker/honker.test.ts @@ -0,0 +1,41 @@ +import chai, { expect } from 'chai'; +import asPromised from 'chai-as-promised'; +import { describe, it } from 'mocha'; +import { transientDb } from '../utils/TransientTestUtils.ts'; +import { Queue, type HonkerJob } from '../../common/database/honker/HonkerQueue.ts'; +import type { JsonPlayObject, QueueContext } from '../../../core/Atomic.ts'; +import { generateJsonPlay } from '../../../core/tests/utils/PlayTestUtils.ts'; +import honker from '@russellthehippo/honker-node'; + +chai.use(asPromised); + + +describe('Expected behavior for queues', function () { + + it('enqueues and claims a job', async function () { + const db = await transientDb(); + const a = db.$client.location() + const hdb = honker.open(a); + type PlayJob = { context: QueueContext, play: JsonPlayObject }; + const queue = new Queue(db, 'source-discovery-1'); + const p = generateJsonPlay(); + queue.enqueue({ context: {}, play: p }); + + const honkerQueue = hdb.queue('source-discovery-1'); + + honkerQueue.sweepExpired + + const waker = honkerQueue.claimWaker(); + while (true) { + const job = await waker.next('worker-1') as unknown as HonkerJob; + if (!job) break; + try { + expect(job.payload.play.data.track).eq(p.data.track); + job.ack(); + break; + } catch (err) { + job.retry(60, String(err)); + } + } + }); +}); \ No newline at end of file