Merged
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
4 changes: 2 additions & 2 deletions package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion package.json
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
{
"name": "@athenna/queue",
"version": "5.32.0",
"version": "5.33.0",
"description": "The Athenna queue handler.",
"license": "MIT",
"author": "João Lenon <lenon@athenna.io>",
Expand Down
1 change: 1 addition & 0 deletions src/index.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -19,6 +19,7 @@ export * from '#src/drivers/DatabaseDriver'
export * from '#src/factories/ConnectionFactory'

export * from '#src/facades/Queue'
export * from '#src/worker/BaseWorker'
export * from '#src/worker/WorkerImpl'
export * from '#src/providers/QueueProvider'
export * from '#src/providers/WorkerProvider'
Expand Down
12 changes: 12 additions & 0 deletions src/types/ConnectionOptions.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -125,5 +125,17 @@ export type ConnectionOptions = {
* @default Parser.timeToMs('5m')
*/
workerTimeoutMs?: number

/**
* Define how many independent consumer loops run in parallel for this
* connection. Each loop still pulls and processes ONE job at a time, so a
* value of `N` yields an effective concurrency of `N`. This is honored by
* the `@athenna/event` consumer; raise it to drain a backed-up queue faster
* without spinning up extra processes. When the option is `null`/unset it
* defaults to `0`, which falls back to a single serial loop (one-by-one).
*
* @default Config.get(`queue.connections.${connection}.workerConcurrency`, 0)
*/
workerConcurrency?: number
}
}
7 changes: 5 additions & 2 deletions src/types/WorkerOptions.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -16,9 +16,12 @@ export type WorkerOptions = {
name?: string

/**
* Define how much instances of the same worker could run in parallel.
* Define how many instances of the same worker run in parallel. Each
* instance still processes one job at a time, so a value of `N` yields an
* effective concurrency of `N`. When omitted, the worker falls back to the
* connection's `workerConcurrency` config and, if that is `0`/unset, to `1`.
*
* @default 1
* @default Config.get(`queue.connections.${connection}.workerConcurrency`, 1)
*/
concurrency?: number

Expand Down
65 changes: 65 additions & 0 deletions src/worker/BaseWorker.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,65 @@
/**
* @athenna/queue
*
* (c) João Lenon <lenon@athenna.io>
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/

import 'reflect-metadata'

import { Queue } from '#src/facades/Queue'
import { Annotation } from '@athenna/ioc'
import type { QueueImpl } from '#src/queue/QueueImpl'

/**
* Base class for workers. Extend it to get a queue instance already bound to
* the worker's own connection, so you don't need to call `Queue.connection()`
* on every operation.
*
* @example
* ```ts
* @Worker()
* export class HelloWorker extends BaseWorker {
* public async handle(ctx: Context) {
* await this.queue.add({ hello: 'world' })
* }
* }
* ```
*/
export class BaseWorker {
/**
* Cached queue instance bound to this worker's connection.
*/
private _queue?: QueueImpl

/**
* The queue connection name of this worker. It is resolved from the
* worker's `@Worker({ connection })` metadata, falling back to the default
* connection (`queue.default`) when the worker is not annotated.
*/
public get connection() {
const meta = Annotation.getMeta(this.constructor)

return meta?.connection ?? Config.get('queue.default')
}

/**
* A queue instance already bound to this worker's connection. Use it to
* enqueue or inspect jobs without calling `Queue.connection(...)` on every
* operation.
*
* @example
* ```ts
* await this.queue.add({ email: 'lenon@athenna.io' })
* ```
*/
public get queue() {
if (!this._queue) {
this._queue = Queue.connection(this.connection)
}

return this._queue
}
}
25 changes: 24 additions & 1 deletion src/worker/WorkerTaskBuilder.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -200,13 +200,36 @@ export class WorkerTaskBuilder {
return
}

const n = this.worker.concurrency ?? 1
const n = this.resolveConcurrency()

for (let i = 0; i < n; i++) {
this.spawn()
}
}

/**
* Resolve how many worker loops to spawn. An explicit `concurrency` set on
* the worker (e.g. via `@Worker({ concurrency })`) wins, then the worker
* `options.workerConcurrency`, then the connection's `workerConcurrency`
* config. When none is a positive number it falls back to a single serial
* loop.
*/
private resolveConcurrency(): number {
const explicit =
this.worker.concurrency ?? this.worker.options?.workerConcurrency

if (Is.Number(explicit) && explicit > 0) {
return explicit
}

const configured = Config.get(
`queue.connections.${this.worker.connection}.workerConcurrency`,
0
)

return configured > 0 ? configured : 1
}

/**
* Use spawn to force a worker instance to run.
*/
Expand Down
4 changes: 2 additions & 2 deletions templates/worker.edge
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
import { Worker, type Context } from '@athenna/queue'
import { Worker, BaseWorker, type Context } from '@athenna/queue'

@Worker()
export class {{ namePascal }} {
export class {{ namePascal }} extends BaseWorker {
public async handle(ctx: Context) {
//
}
Expand Down
7 changes: 7 additions & 0 deletions tests/fixtures/config/queue.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -61,6 +61,13 @@ export default {
workerTimeoutMs: 200
},

memoryConcurrent: {
driver: 'memory',
queue: 'default',
deadletter: 'deadletter',
workerConcurrency: 3
},

aws_sqs: {
driver: 'aws_sqs',
type: 'standard',
Expand Down
82 changes: 82 additions & 0 deletions tests/unit/worker/BaseWorkerTest.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,82 @@
/**
* @athenna/queue
*
* (c) João Lenon <lenon@athenna.io>
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/

import { Path } from '@athenna/common'
import { Worker, BaseWorker, QueueImpl } from '#src'
import { LoggerProvider } from '@athenna/logger'
import { QueueProvider } from '#src/providers/QueueProvider'
import { Test, BeforeEach, AfterEach, type Context } from '@athenna/test'

@Worker({ connection: 'fake' })
class FakeConnectionWorker extends BaseWorker {
public async handle() {}
}

@Worker({ connection: 'memory' })
class MemoryConnectionWorker extends BaseWorker {
public async handle() {}
}

@Worker()
class DefaultConnectionWorker extends BaseWorker {
public async handle() {}
}

export class BaseWorkerTest {
@BeforeEach()
public async beforeEach() {
await Config.loadAll(Path.fixtures('config'))

new LoggerProvider().register()
new QueueProvider().register()
}

@AfterEach()
public async afterEach() {
await new QueueProvider().shutdown()

ioc.reconstruct()
Config.clear()
}

@Test()
public async shouldResolveTheConnectionFromTheWorkerMetadata({ assert }: Context) {
const worker = new FakeConnectionWorker()

assert.equal(worker.connection, 'fake')
assert.equal(worker.queue.connectionName, 'fake')
}

@Test()
public async shouldFallBackToTheDefaultConnectionWhenNotAnnotatedWithOne({ assert }: Context) {
const worker = new DefaultConnectionWorker()

assert.equal(worker.connection, Config.get('queue.default'))
assert.equal(worker.queue.connectionName, Config.get('queue.default'))
}

@Test()
public async shouldExposeAReadyToUseQueueInstanceBoundToTheWorkerConnection({ assert }: Context) {
const worker = new MemoryConnectionWorker()

assert.instanceOf(worker.queue, QueueImpl)
assert.equal(worker.queue.connectionName, 'memory')
assert.isTrue(worker.queue.isConnected())
}

@Test()
public async shouldCacheTheQueueInstanceBetweenAccesses({ assert }: Context) {
const worker = new FakeConnectionWorker()

const first = worker.queue
const second = worker.queue

assert.isTrue(first === second)
}
}
62 changes: 62 additions & 0 deletions tests/unit/worker/WorkerImplTest.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -176,6 +176,68 @@ export class WorkerImplTest {
assert.deepEqual(task?.worker.connection, 'fake')
}

@Test()
public async shouldSpawnASingleLoopWhenNoConcurrencyIsConfigured({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('default_concurrency')
.connection('memory')
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 1)

builder.stop()
}

@Test()
public async shouldSpawnConcurrentLoopsBasedOnConnectionWorkerConcurrencyConfig({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('config_concurrency')
.connection('memoryConcurrent')
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 3)

builder.stop()
}

@Test()
public async shouldSpawnConcurrentLoopsFromWorkerOptionsWorkerConcurrency({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('options_concurrency')
.connection('memory')
.options({ workerConcurrency: 4 })
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 4)

builder.stop()
}

@Test()
public async shouldLetExplicitConcurrencyOverrideTheConnectionWorkerConcurrencyConfig({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('explicit_concurrency')
.connection('memoryConcurrent')
.concurrency(5)
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 5)

builder.stop()
}

@Test()
public async shouldBeAbleToCreateAWorkerTaskWithCustomOptions({ assert }: Context) {
Queue.worker()
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content
Merged
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
4 changes: 2 additions & 2 deletions package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion package.json
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
{
"name": "@athenna/queue",
"version": "5.32.0",
"version": "5.33.0",
"description": "The Athenna queue handler.",
"license": "MIT",
"author": "João Lenon <lenon@athenna.io>",
Expand Down
1 change: 1 addition & 0 deletions src/index.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -19,6 +19,7 @@ export * from '#src/drivers/DatabaseDriver'
export * from '#src/factories/ConnectionFactory'

export * from '#src/facades/Queue'
export * from '#src/worker/BaseWorker'
export * from '#src/worker/WorkerImpl'
export * from '#src/providers/QueueProvider'
export * from '#src/providers/WorkerProvider'
Expand Down
12 changes: 12 additions & 0 deletions src/types/ConnectionOptions.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -125,5 +125,17 @@ export type ConnectionOptions = {
* @default Parser.timeToMs('5m')
*/
workerTimeoutMs?: number

/**
* Define how many independent consumer loops run in parallel for this
* connection. Each loop still pulls and processes ONE job at a time, so a
* value of `N` yields an effective concurrency of `N`. This is honored by
* the `@athenna/event` consumer; raise it to drain a backed-up queue faster
* without spinning up extra processes. When the option is `null`/unset it
* defaults to `0`, which falls back to a single serial loop (one-by-one).
*
* @default Config.get(`queue.connections.${connection}.workerConcurrency`, 0)
*/
workerConcurrency?: number
}
}
7 changes: 5 additions & 2 deletions src/types/WorkerOptions.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -16,9 +16,12 @@ export type WorkerOptions = {
name?: string

/**
* Define how much instances of the same worker could run in parallel.
* Define how many instances of the same worker run in parallel. Each
* instance still processes one job at a time, so a value of `N` yields an
* effective concurrency of `N`. When omitted, the worker falls back to the
* connection's `workerConcurrency` config and, if that is `0`/unset, to `1`.
*
* @default 1
* @default Config.get(`queue.connections.${connection}.workerConcurrency`, 1)
*/
concurrency?: number

Expand Down
65 changes: 65 additions & 0 deletions src/worker/BaseWorker.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,65 @@
/**
* @athenna/queue
*
* (c) João Lenon <lenon@athenna.io>
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/

import 'reflect-metadata'

import { Queue } from '#src/facades/Queue'
import { Annotation } from '@athenna/ioc'
import type { QueueImpl } from '#src/queue/QueueImpl'

/**
* Base class for workers. Extend it to get a queue instance already bound to
* the worker's own connection, so you don't need to call `Queue.connection()`
* on every operation.
*
* @example
* ```ts
* @Worker()
* export class HelloWorker extends BaseWorker {
* public async handle(ctx: Context) {
* await this.queue.add({ hello: 'world' })
* }
* }
* ```
*/
export class BaseWorker {
/**
* Cached queue instance bound to this worker's connection.
*/
private _queue?: QueueImpl

/**
* The queue connection name of this worker. It is resolved from the
* worker's `@Worker({ connection })` metadata, falling back to the default
* connection (`queue.default`) when the worker is not annotated.
*/
public get connection() {
const meta = Annotation.getMeta(this.constructor)

return meta?.connection ?? Config.get('queue.default')
}

/**
* A queue instance already bound to this worker's connection. Use it to
* enqueue or inspect jobs without calling `Queue.connection(...)` on every
* operation.
*
* @example
* ```ts
* await this.queue.add({ email: 'lenon@athenna.io' })
* ```
*/
public get queue() {
if (!this._queue) {
this._queue = Queue.connection(this.connection)
}

return this._queue
}
}
25 changes: 24 additions & 1 deletion src/worker/WorkerTaskBuilder.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -200,13 +200,36 @@ export class WorkerTaskBuilder {
return
}

const n = this.worker.concurrency ?? 1
const n = this.resolveConcurrency()

for (let i = 0; i < n; i++) {
this.spawn()
}
}

/**
* Resolve how many worker loops to spawn. An explicit `concurrency` set on
* the worker (e.g. via `@Worker({ concurrency })`) wins, then the worker
* `options.workerConcurrency`, then the connection's `workerConcurrency`
* config. When none is a positive number it falls back to a single serial
* loop.
*/
private resolveConcurrency(): number {
const explicit =
this.worker.concurrency ?? this.worker.options?.workerConcurrency

if (Is.Number(explicit) && explicit > 0) {
return explicit
}

const configured = Config.get(
`queue.connections.${this.worker.connection}.workerConcurrency`,
0
)

return configured > 0 ? configured : 1
}

/**
* Use spawn to force a worker instance to run.
*/
Expand Down
4 changes: 2 additions & 2 deletions templates/worker.edge
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
import { Worker, type Context } from '@athenna/queue'
import { Worker, BaseWorker, type Context } from '@athenna/queue'

@Worker()
export class {{ namePascal }} {
export class {{ namePascal }} extends BaseWorker {
public async handle(ctx: Context) {
//
}
Expand Down
7 changes: 7 additions & 0 deletions tests/fixtures/config/queue.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -61,6 +61,13 @@ export default {
workerTimeoutMs: 200
},

memoryConcurrent: {
driver: 'memory',
queue: 'default',
deadletter: 'deadletter',
workerConcurrency: 3
},

aws_sqs: {
driver: 'aws_sqs',
type: 'standard',
Expand Down
82 changes: 82 additions & 0 deletions tests/unit/worker/BaseWorkerTest.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,82 @@
/**
* @athenna/queue
*
* (c) João Lenon <lenon@athenna.io>
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/

import { Path } from '@athenna/common'
import { Worker, BaseWorker, QueueImpl } from '#src'
import { LoggerProvider } from '@athenna/logger'
import { QueueProvider } from '#src/providers/QueueProvider'
import { Test, BeforeEach, AfterEach, type Context } from '@athenna/test'

@Worker({ connection: 'fake' })
class FakeConnectionWorker extends BaseWorker {
public async handle() {}
}

@Worker({ connection: 'memory' })
class MemoryConnectionWorker extends BaseWorker {
public async handle() {}
}

@Worker()
class DefaultConnectionWorker extends BaseWorker {
public async handle() {}
}

export class BaseWorkerTest {
@BeforeEach()
public async beforeEach() {
await Config.loadAll(Path.fixtures('config'))

new LoggerProvider().register()
new QueueProvider().register()
}

@AfterEach()
public async afterEach() {
await new QueueProvider().shutdown()

ioc.reconstruct()
Config.clear()
}

@Test()
public async shouldResolveTheConnectionFromTheWorkerMetadata({ assert }: Context) {
const worker = new FakeConnectionWorker()

assert.equal(worker.connection, 'fake')
assert.equal(worker.queue.connectionName, 'fake')
}

@Test()
public async shouldFallBackToTheDefaultConnectionWhenNotAnnotatedWithOne({ assert }: Context) {
const worker = new DefaultConnectionWorker()

assert.equal(worker.connection, Config.get('queue.default'))
assert.equal(worker.queue.connectionName, Config.get('queue.default'))
}

@Test()
public async shouldExposeAReadyToUseQueueInstanceBoundToTheWorkerConnection({ assert }: Context) {
const worker = new MemoryConnectionWorker()

assert.instanceOf(worker.queue, QueueImpl)
assert.equal(worker.queue.connectionName, 'memory')
assert.isTrue(worker.queue.isConnected())
}

@Test()
public async shouldCacheTheQueueInstanceBetweenAccesses({ assert }: Context) {
const worker = new FakeConnectionWorker()

const first = worker.queue
const second = worker.queue

assert.isTrue(first === second)
}
}
62 changes: 62 additions & 0 deletions tests/unit/worker/WorkerImplTest.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -176,6 +176,68 @@ export class WorkerImplTest {
assert.deepEqual(task?.worker.connection, 'fake')
}

@Test()
public async shouldSpawnASingleLoopWhenNoConcurrencyIsConfigured({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('default_concurrency')
.connection('memory')
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 1)

builder.stop()
}

@Test()
public async shouldSpawnConcurrentLoopsBasedOnConnectionWorkerConcurrencyConfig({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('config_concurrency')
.connection('memoryConcurrent')
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 3)

builder.stop()
}

@Test()
public async shouldSpawnConcurrentLoopsFromWorkerOptionsWorkerConcurrency({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('options_concurrency')
.connection('memory')
.options({ workerConcurrency: 4 })
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 4)

builder.stop()
}

@Test()
public async shouldLetExplicitConcurrencyOverrideTheConnectionWorkerConcurrencyConfig({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('explicit_concurrency')
.connection('memoryConcurrent')
.concurrency(5)
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 5)

builder.stop()
}

@Test()
public async shouldBeAbleToCreateAWorkerTaskWithCustomOptions({ assert }: Context) {
Queue.worker()
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
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
4 changes: 2 additions & 2 deletions package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion package.json
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
{
"name": "@athenna/queue",
"version": "5.32.0",
"version": "5.33.0",
"description": "The Athenna queue handler.",
"license": "MIT",
"author": "João Lenon <lenon@athenna.io>",
Expand Down
1 change: 1 addition & 0 deletions src/index.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -19,6 +19,7 @@ export * from '#src/drivers/DatabaseDriver'
export * from '#src/factories/ConnectionFactory'

export * from '#src/facades/Queue'
export * from '#src/worker/BaseWorker'
export * from '#src/worker/WorkerImpl'
export * from '#src/providers/QueueProvider'
export * from '#src/providers/WorkerProvider'
Expand Down
12 changes: 12 additions & 0 deletions src/types/ConnectionOptions.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -125,5 +125,17 @@ export type ConnectionOptions = {
* @default Parser.timeToMs('5m')
*/
workerTimeoutMs?: number

/**
* Define how many independent consumer loops run in parallel for this
* connection. Each loop still pulls and processes ONE job at a time, so a
* value of `N` yields an effective concurrency of `N`. This is honored by
* the `@athenna/event` consumer; raise it to drain a backed-up queue faster
* without spinning up extra processes. When the option is `null`/unset it
* defaults to `0`, which falls back to a single serial loop (one-by-one).
*
* @default Config.get(`queue.connections.${connection}.workerConcurrency`, 0)
*/
workerConcurrency?: number
}
}
7 changes: 5 additions & 2 deletions src/types/WorkerOptions.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -16,9 +16,12 @@ export type WorkerOptions = {
name?: string

/**
* Define how much instances of the same worker could run in parallel.
* Define how many instances of the same worker run in parallel. Each
* instance still processes one job at a time, so a value of `N` yields an
* effective concurrency of `N`. When omitted, the worker falls back to the
* connection's `workerConcurrency` config and, if that is `0`/unset, to `1`.
*
* @default 1
* @default Config.get(`queue.connections.${connection}.workerConcurrency`, 1)
*/
concurrency?: number

Expand Down
65 changes: 65 additions & 0 deletions src/worker/BaseWorker.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,65 @@
/**
* @athenna/queue
*
* (c) João Lenon <lenon@athenna.io>
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/

import 'reflect-metadata'

import { Queue } from '#src/facades/Queue'
import { Annotation } from '@athenna/ioc'
import type { QueueImpl } from '#src/queue/QueueImpl'

/**
* Base class for workers. Extend it to get a queue instance already bound to
* the worker's own connection, so you don't need to call `Queue.connection()`
* on every operation.
*
* @example
* ```ts
* @Worker()
* export class HelloWorker extends BaseWorker {
* public async handle(ctx: Context) {
* await this.queue.add({ hello: 'world' })
* }
* }
* ```
*/
export class BaseWorker {
/**
* Cached queue instance bound to this worker's connection.
*/
private _queue?: QueueImpl

/**
* The queue connection name of this worker. It is resolved from the
* worker's `@Worker({ connection })` metadata, falling back to the default
* connection (`queue.default`) when the worker is not annotated.
*/
public get connection() {
const meta = Annotation.getMeta(this.constructor)

return meta?.connection ?? Config.get('queue.default')
}

/**
* A queue instance already bound to this worker's connection. Use it to
* enqueue or inspect jobs without calling `Queue.connection(...)` on every
* operation.
*
* @example
* ```ts
* await this.queue.add({ email: 'lenon@athenna.io' })
* ```
*/
public get queue() {
if (!this._queue) {
this._queue = Queue.connection(this.connection)
}

return this._queue
}
}
25 changes: 24 additions & 1 deletion src/worker/WorkerTaskBuilder.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -200,13 +200,36 @@ export class WorkerTaskBuilder {
return
}

const n = this.worker.concurrency ?? 1
const n = this.resolveConcurrency()

for (let i = 0; i < n; i++) {
this.spawn()
}
}

/**
* Resolve how many worker loops to spawn. An explicit `concurrency` set on
* the worker (e.g. via `@Worker({ concurrency })`) wins, then the worker
* `options.workerConcurrency`, then the connection's `workerConcurrency`
* config. When none is a positive number it falls back to a single serial
* loop.
*/
private resolveConcurrency(): number {
const explicit =
this.worker.concurrency ?? this.worker.options?.workerConcurrency

if (Is.Number(explicit) && explicit > 0) {
return explicit
}

const configured = Config.get(
`queue.connections.${this.worker.connection}.workerConcurrency`,
0
)

return configured > 0 ? configured : 1
}

/**
* Use spawn to force a worker instance to run.
*/
Expand Down
4 changes: 2 additions & 2 deletions templates/worker.edge
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
import { Worker, type Context } from '@athenna/queue'
import { Worker, BaseWorker, type Context } from '@athenna/queue'

@Worker()
export class {{ namePascal }} {
export class {{ namePascal }} extends BaseWorker {
public async handle(ctx: Context) {
//
}
Expand Down
7 changes: 7 additions & 0 deletions tests/fixtures/config/queue.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -61,6 +61,13 @@ export default {
workerTimeoutMs: 200
},

memoryConcurrent: {
driver: 'memory',
queue: 'default',
deadletter: 'deadletter',
workerConcurrency: 3
},

aws_sqs: {
driver: 'aws_sqs',
type: 'standard',
Expand Down
82 changes: 82 additions & 0 deletions tests/unit/worker/BaseWorkerTest.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,82 @@
/**
* @athenna/queue
*
* (c) João Lenon <lenon@athenna.io>
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/

import { Path } from '@athenna/common'
import { Worker, BaseWorker, QueueImpl } from '#src'
import { LoggerProvider } from '@athenna/logger'
import { QueueProvider } from '#src/providers/QueueProvider'
import { Test, BeforeEach, AfterEach, type Context } from '@athenna/test'

@Worker({ connection: 'fake' })
class FakeConnectionWorker extends BaseWorker {
public async handle() {}
}

@Worker({ connection: 'memory' })
class MemoryConnectionWorker extends BaseWorker {
public async handle() {}
}

@Worker()
class DefaultConnectionWorker extends BaseWorker {
public async handle() {}
}

export class BaseWorkerTest {
@BeforeEach()
public async beforeEach() {
await Config.loadAll(Path.fixtures('config'))

new LoggerProvider().register()
new QueueProvider().register()
}

@AfterEach()
public async afterEach() {
await new QueueProvider().shutdown()

ioc.reconstruct()
Config.clear()
}

@Test()
public async shouldResolveTheConnectionFromTheWorkerMetadata({ assert }: Context) {
const worker = new FakeConnectionWorker()

assert.equal(worker.connection, 'fake')
assert.equal(worker.queue.connectionName, 'fake')
}

@Test()
public async shouldFallBackToTheDefaultConnectionWhenNotAnnotatedWithOne({ assert }: Context) {
const worker = new DefaultConnectionWorker()

assert.equal(worker.connection, Config.get('queue.default'))
assert.equal(worker.queue.connectionName, Config.get('queue.default'))
}

@Test()
public async shouldExposeAReadyToUseQueueInstanceBoundToTheWorkerConnection({ assert }: Context) {
const worker = new MemoryConnectionWorker()

assert.instanceOf(worker.queue, QueueImpl)
assert.equal(worker.queue.connectionName, 'memory')
assert.isTrue(worker.queue.isConnected())
}

@Test()
public async shouldCacheTheQueueInstanceBetweenAccesses({ assert }: Context) {
const worker = new FakeConnectionWorker()

const first = worker.queue
const second = worker.queue

assert.isTrue(first === second)
}
}
62 changes: 62 additions & 0 deletions tests/unit/worker/WorkerImplTest.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -176,6 +176,68 @@ export class WorkerImplTest {
assert.deepEqual(task?.worker.connection, 'fake')
}

@Test()
public async shouldSpawnASingleLoopWhenNoConcurrencyIsConfigured({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('default_concurrency')
.connection('memory')
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 1)

builder.stop()
}

@Test()
public async shouldSpawnConcurrentLoopsBasedOnConnectionWorkerConcurrencyConfig({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('config_concurrency')
.connection('memoryConcurrent')
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 3)

builder.stop()
}

@Test()
public async shouldSpawnConcurrentLoopsFromWorkerOptionsWorkerConcurrency({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('options_concurrency')
.connection('memory')
.options({ workerConcurrency: 4 })
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 4)

builder.stop()
}

@Test()
public async shouldLetExplicitConcurrencyOverrideTheConnectionWorkerConcurrencyConfig({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('explicit_concurrency')
.connection('memoryConcurrent')
.concurrency(5)
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 5)

builder.stop()
}

@Test()
public async shouldBeAbleToCreateAWorkerTaskWithCustomOptions({ assert }: Context) {
Queue.worker()
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
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
4 changes: 2 additions & 2 deletions package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion package.json
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
{
"name": "@athenna/queue",
"version": "5.32.0",
"version": "5.33.0",
"description": "The Athenna queue handler.",
"license": "MIT",
"author": "João Lenon <lenon@athenna.io>",
Expand Down
1 change: 1 addition & 0 deletions src/index.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -19,6 +19,7 @@ export * from '#src/drivers/DatabaseDriver'
export * from '#src/factories/ConnectionFactory'

export * from '#src/facades/Queue'
export * from '#src/worker/BaseWorker'
export * from '#src/worker/WorkerImpl'
export * from '#src/providers/QueueProvider'
export * from '#src/providers/WorkerProvider'
Expand Down
12 changes: 12 additions & 0 deletions src/types/ConnectionOptions.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -125,5 +125,17 @@ export type ConnectionOptions = {
* @default Parser.timeToMs('5m')
*/
workerTimeoutMs?: number

/**
* Define how many independent consumer loops run in parallel for this
* connection. Each loop still pulls and processes ONE job at a time, so a
* value of `N` yields an effective concurrency of `N`. This is honored by
* the `@athenna/event` consumer; raise it to drain a backed-up queue faster
* without spinning up extra processes. When the option is `null`/unset it
* defaults to `0`, which falls back to a single serial loop (one-by-one).
*
* @default Config.get(`queue.connections.${connection}.workerConcurrency`, 0)
*/
workerConcurrency?: number
}
}
7 changes: 5 additions & 2 deletions src/types/WorkerOptions.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -16,9 +16,12 @@ export type WorkerOptions = {
name?: string

/**
* Define how much instances of the same worker could run in parallel.
* Define how many instances of the same worker run in parallel. Each
* instance still processes one job at a time, so a value of `N` yields an
* effective concurrency of `N`. When omitted, the worker falls back to the
* connection's `workerConcurrency` config and, if that is `0`/unset, to `1`.
*
* @default 1
* @default Config.get(`queue.connections.${connection}.workerConcurrency`, 1)
*/
concurrency?: number

Expand Down
65 changes: 65 additions & 0 deletions src/worker/BaseWorker.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,65 @@
/**
* @athenna/queue
*
* (c) João Lenon <lenon@athenna.io>
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/

import 'reflect-metadata'

import { Queue } from '#src/facades/Queue'
import { Annotation } from '@athenna/ioc'
import type { QueueImpl } from '#src/queue/QueueImpl'

/**
* Base class for workers. Extend it to get a queue instance already bound to
* the worker's own connection, so you don't need to call `Queue.connection()`
* on every operation.
*
* @example
* ```ts
* @Worker()
* export class HelloWorker extends BaseWorker {
* public async handle(ctx: Context) {
* await this.queue.add({ hello: 'world' })
* }
* }
* ```
*/
export class BaseWorker {
/**
* Cached queue instance bound to this worker's connection.
*/
private _queue?: QueueImpl

/**
* The queue connection name of this worker. It is resolved from the
* worker's `@Worker({ connection })` metadata, falling back to the default
* connection (`queue.default`) when the worker is not annotated.
*/
public get connection() {
const meta = Annotation.getMeta(this.constructor)

return meta?.connection ?? Config.get('queue.default')
}

/**
* A queue instance already bound to this worker's connection. Use it to
* enqueue or inspect jobs without calling `Queue.connection(...)` on every
* operation.
*
* @example
* ```ts
* await this.queue.add({ email: 'lenon@athenna.io' })
* ```
*/
public get queue() {
if (!this._queue) {
this._queue = Queue.connection(this.connection)
}

return this._queue
}
}
25 changes: 24 additions & 1 deletion src/worker/WorkerTaskBuilder.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -200,13 +200,36 @@ export class WorkerTaskBuilder {
return
}

const n = this.worker.concurrency ?? 1
const n = this.resolveConcurrency()

for (let i = 0; i < n; i++) {
this.spawn()
}
}

/**
* Resolve how many worker loops to spawn. An explicit `concurrency` set on
* the worker (e.g. via `@Worker({ concurrency })`) wins, then the worker
* `options.workerConcurrency`, then the connection's `workerConcurrency`
* config. When none is a positive number it falls back to a single serial
* loop.
*/
private resolveConcurrency(): number {
const explicit =
this.worker.concurrency ?? this.worker.options?.workerConcurrency

if (Is.Number(explicit) && explicit > 0) {
return explicit
}

const configured = Config.get(
`queue.connections.${this.worker.connection}.workerConcurrency`,
0
)

return configured > 0 ? configured : 1
}

/**
* Use spawn to force a worker instance to run.
*/
Expand Down
4 changes: 2 additions & 2 deletions templates/worker.edge
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
import { Worker, type Context } from '@athenna/queue'
import { Worker, BaseWorker, type Context } from '@athenna/queue'

@Worker()
export class {{ namePascal }} {
export class {{ namePascal }} extends BaseWorker {
public async handle(ctx: Context) {
//
}
Expand Down
7 changes: 7 additions & 0 deletions tests/fixtures/config/queue.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -61,6 +61,13 @@ export default {
workerTimeoutMs: 200
},

memoryConcurrent: {
driver: 'memory',
queue: 'default',
deadletter: 'deadletter',
workerConcurrency: 3
},

aws_sqs: {
driver: 'aws_sqs',
type: 'standard',
Expand Down
82 changes: 82 additions & 0 deletions tests/unit/worker/BaseWorkerTest.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,82 @@
/**
* @athenna/queue
*
* (c) João Lenon <lenon@athenna.io>
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/

import { Path } from '@athenna/common'
import { Worker, BaseWorker, QueueImpl } from '#src'
import { LoggerProvider } from '@athenna/logger'
import { QueueProvider } from '#src/providers/QueueProvider'
import { Test, BeforeEach, AfterEach, type Context } from '@athenna/test'

@Worker({ connection: 'fake' })
class FakeConnectionWorker extends BaseWorker {
public async handle() {}
}

@Worker({ connection: 'memory' })
class MemoryConnectionWorker extends BaseWorker {
public async handle() {}
}

@Worker()
class DefaultConnectionWorker extends BaseWorker {
public async handle() {}
}

export class BaseWorkerTest {
@BeforeEach()
public async beforeEach() {
await Config.loadAll(Path.fixtures('config'))

new LoggerProvider().register()
new QueueProvider().register()
}

@AfterEach()
public async afterEach() {
await new QueueProvider().shutdown()

ioc.reconstruct()
Config.clear()
}

@Test()
public async shouldResolveTheConnectionFromTheWorkerMetadata({ assert }: Context) {
const worker = new FakeConnectionWorker()

assert.equal(worker.connection, 'fake')
assert.equal(worker.queue.connectionName, 'fake')
}

@Test()
public async shouldFallBackToTheDefaultConnectionWhenNotAnnotatedWithOne({ assert }: Context) {
const worker = new DefaultConnectionWorker()

assert.equal(worker.connection, Config.get('queue.default'))
assert.equal(worker.queue.connectionName, Config.get('queue.default'))
}

@Test()
public async shouldExposeAReadyToUseQueueInstanceBoundToTheWorkerConnection({ assert }: Context) {
const worker = new MemoryConnectionWorker()

assert.instanceOf(worker.queue, QueueImpl)
assert.equal(worker.queue.connectionName, 'memory')
assert.isTrue(worker.queue.isConnected())
}

@Test()
public async shouldCacheTheQueueInstanceBetweenAccesses({ assert }: Context) {
const worker = new FakeConnectionWorker()

const first = worker.queue
const second = worker.queue

assert.isTrue(first === second)
}
}
62 changes: 62 additions & 0 deletions tests/unit/worker/WorkerImplTest.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -176,6 +176,68 @@ export class WorkerImplTest {
assert.deepEqual(task?.worker.connection, 'fake')
}

@Test()
public async shouldSpawnASingleLoopWhenNoConcurrencyIsConfigured({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('default_concurrency')
.connection('memory')
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 1)

builder.stop()
}

@Test()
public async shouldSpawnConcurrentLoopsBasedOnConnectionWorkerConcurrencyConfig({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('config_concurrency')
.connection('memoryConcurrent')
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 3)

builder.stop()
}

@Test()
public async shouldSpawnConcurrentLoopsFromWorkerOptionsWorkerConcurrency({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('options_concurrency')
.connection('memory')
.options({ workerConcurrency: 4 })
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 4)

builder.stop()
}

@Test()
public async shouldLetExplicitConcurrencyOverrideTheConnectionWorkerConcurrencyConfig({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('explicit_concurrency')
.connection('memoryConcurrent')
.concurrency(5)
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 5)

builder.stop()
}

@Test()
public async shouldBeAbleToCreateAWorkerTaskWithCustomOptions({ assert }: Context) {
Queue.worker()
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content
Merged
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
4 changes: 2 additions & 2 deletions package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion package.json
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
{
"name": "@athenna/queue",
"version": "5.32.0",
"version": "5.33.0",
"description": "The Athenna queue handler.",
"license": "MIT",
"author": "João Lenon <lenon@athenna.io>",
Expand Down
1 change: 1 addition & 0 deletions src/index.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -19,6 +19,7 @@ export * from '#src/drivers/DatabaseDriver'
export * from '#src/factories/ConnectionFactory'

export * from '#src/facades/Queue'
export * from '#src/worker/BaseWorker'
export * from '#src/worker/WorkerImpl'
export * from '#src/providers/QueueProvider'
export * from '#src/providers/WorkerProvider'
Expand Down
12 changes: 12 additions & 0 deletions src/types/ConnectionOptions.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -125,5 +125,17 @@ export type ConnectionOptions = {
* @default Parser.timeToMs('5m')
*/
workerTimeoutMs?: number

/**
* Define how many independent consumer loops run in parallel for this
* connection. Each loop still pulls and processes ONE job at a time, so a
* value of `N` yields an effective concurrency of `N`. This is honored by
* the `@athenna/event` consumer; raise it to drain a backed-up queue faster
* without spinning up extra processes. When the option is `null`/unset it
* defaults to `0`, which falls back to a single serial loop (one-by-one).
*
* @default Config.get(`queue.connections.${connection}.workerConcurrency`, 0)
*/
workerConcurrency?: number
}
}
7 changes: 5 additions & 2 deletions src/types/WorkerOptions.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -16,9 +16,12 @@ export type WorkerOptions = {
name?: string

/**
* Define how much instances of the same worker could run in parallel.
* Define how many instances of the same worker run in parallel. Each
* instance still processes one job at a time, so a value of `N` yields an
* effective concurrency of `N`. When omitted, the worker falls back to the
* connection's `workerConcurrency` config and, if that is `0`/unset, to `1`.
*
* @default 1
* @default Config.get(`queue.connections.${connection}.workerConcurrency`, 1)
*/
concurrency?: number

Expand Down
65 changes: 65 additions & 0 deletions src/worker/BaseWorker.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,65 @@
/**
* @athenna/queue
*
* (c) João Lenon <lenon@athenna.io>
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/

import 'reflect-metadata'

import { Queue } from '#src/facades/Queue'
import { Annotation } from '@athenna/ioc'
import type { QueueImpl } from '#src/queue/QueueImpl'

/**
* Base class for workers. Extend it to get a queue instance already bound to
* the worker's own connection, so you don't need to call `Queue.connection()`
* on every operation.
*
* @example
* ```ts
* @Worker()
* export class HelloWorker extends BaseWorker {
* public async handle(ctx: Context) {
* await this.queue.add({ hello: 'world' })
* }
* }
* ```
*/
export class BaseWorker {
/**
* Cached queue instance bound to this worker's connection.
*/
private _queue?: QueueImpl

/**
* The queue connection name of this worker. It is resolved from the
* worker's `@Worker({ connection })` metadata, falling back to the default
* connection (`queue.default`) when the worker is not annotated.
*/
public get connection() {
const meta = Annotation.getMeta(this.constructor)

return meta?.connection ?? Config.get('queue.default')
}

/**
* A queue instance already bound to this worker's connection. Use it to
* enqueue or inspect jobs without calling `Queue.connection(...)` on every
* operation.
*
* @example
* ```ts
* await this.queue.add({ email: 'lenon@athenna.io' })
* ```
*/
public get queue() {
if (!this._queue) {
this._queue = Queue.connection(this.connection)
}

return this._queue
}
}
25 changes: 24 additions & 1 deletion src/worker/WorkerTaskBuilder.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -200,13 +200,36 @@ export class WorkerTaskBuilder {
return
}

const n = this.worker.concurrency ?? 1
const n = this.resolveConcurrency()

for (let i = 0; i < n; i++) {
this.spawn()
}
}

/**
* Resolve how many worker loops to spawn. An explicit `concurrency` set on
* the worker (e.g. via `@Worker({ concurrency })`) wins, then the worker
* `options.workerConcurrency`, then the connection's `workerConcurrency`
* config. When none is a positive number it falls back to a single serial
* loop.
*/
private resolveConcurrency(): number {
const explicit =
this.worker.concurrency ?? this.worker.options?.workerConcurrency

if (Is.Number(explicit) && explicit > 0) {
return explicit
}

const configured = Config.get(
`queue.connections.${this.worker.connection}.workerConcurrency`,
0
)

return configured > 0 ? configured : 1
}

/**
* Use spawn to force a worker instance to run.
*/
Expand Down
4 changes: 2 additions & 2 deletions templates/worker.edge
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
import { Worker, type Context } from '@athenna/queue'
import { Worker, BaseWorker, type Context } from '@athenna/queue'

@Worker()
export class {{ namePascal }} {
export class {{ namePascal }} extends BaseWorker {
public async handle(ctx: Context) {
//
}
Expand Down
7 changes: 7 additions & 0 deletions tests/fixtures/config/queue.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -61,6 +61,13 @@ export default {
workerTimeoutMs: 200
},

memoryConcurrent: {
driver: 'memory',
queue: 'default',
deadletter: 'deadletter',
workerConcurrency: 3
},

aws_sqs: {
driver: 'aws_sqs',
type: 'standard',
Expand Down
82 changes: 82 additions & 0 deletions tests/unit/worker/BaseWorkerTest.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,82 @@
/**
* @athenna/queue
*
* (c) João Lenon <lenon@athenna.io>
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/

import { Path } from '@athenna/common'
import { Worker, BaseWorker, QueueImpl } from '#src'
import { LoggerProvider } from '@athenna/logger'
import { QueueProvider } from '#src/providers/QueueProvider'
import { Test, BeforeEach, AfterEach, type Context } from '@athenna/test'

@Worker({ connection: 'fake' })
class FakeConnectionWorker extends BaseWorker {
public async handle() {}
}

@Worker({ connection: 'memory' })
class MemoryConnectionWorker extends BaseWorker {
public async handle() {}
}

@Worker()
class DefaultConnectionWorker extends BaseWorker {
public async handle() {}
}

export class BaseWorkerTest {
@BeforeEach()
public async beforeEach() {
await Config.loadAll(Path.fixtures('config'))

new LoggerProvider().register()
new QueueProvider().register()
}

@AfterEach()
public async afterEach() {
await new QueueProvider().shutdown()

ioc.reconstruct()
Config.clear()
}

@Test()
public async shouldResolveTheConnectionFromTheWorkerMetadata({ assert }: Context) {
const worker = new FakeConnectionWorker()

assert.equal(worker.connection, 'fake')
assert.equal(worker.queue.connectionName, 'fake')
}

@Test()
public async shouldFallBackToTheDefaultConnectionWhenNotAnnotatedWithOne({ assert }: Context) {
const worker = new DefaultConnectionWorker()

assert.equal(worker.connection, Config.get('queue.default'))
assert.equal(worker.queue.connectionName, Config.get('queue.default'))
}

@Test()
public async shouldExposeAReadyToUseQueueInstanceBoundToTheWorkerConnection({ assert }: Context) {
const worker = new MemoryConnectionWorker()

assert.instanceOf(worker.queue, QueueImpl)
assert.equal(worker.queue.connectionName, 'memory')
assert.isTrue(worker.queue.isConnected())
}

@Test()
public async shouldCacheTheQueueInstanceBetweenAccesses({ assert }: Context) {
const worker = new FakeConnectionWorker()

const first = worker.queue
const second = worker.queue

assert.isTrue(first === second)
}
}
62 changes: 62 additions & 0 deletions tests/unit/worker/WorkerImplTest.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -176,6 +176,68 @@ export class WorkerImplTest {
assert.deepEqual(task?.worker.connection, 'fake')
}

@Test()
public async shouldSpawnASingleLoopWhenNoConcurrencyIsConfigured({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('default_concurrency')
.connection('memory')
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 1)

builder.stop()
}

@Test()
public async shouldSpawnConcurrentLoopsBasedOnConnectionWorkerConcurrencyConfig({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('config_concurrency')
.connection('memoryConcurrent')
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 3)

builder.stop()
}

@Test()
public async shouldSpawnConcurrentLoopsFromWorkerOptionsWorkerConcurrency({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('options_concurrency')
.connection('memory')
.options({ workerConcurrency: 4 })
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 4)

builder.stop()
}

@Test()
public async shouldLetExplicitConcurrencyOverrideTheConnectionWorkerConcurrencyConfig({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('explicit_concurrency')
.connection('memoryConcurrent')
.concurrency(5)
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 5)

builder.stop()
}

@Test()
public async shouldBeAbleToCreateAWorkerTaskWithCustomOptions({ assert }: Context) {
Queue.worker()
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
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
4 changes: 2 additions & 2 deletions package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion package.json
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
{
"name": "@athenna/queue",
"version": "5.32.0",
"version": "5.33.0",
"description": "The Athenna queue handler.",
"license": "MIT",
"author": "João Lenon <lenon@athenna.io>",
Expand Down
1 change: 1 addition & 0 deletions src/index.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -19,6 +19,7 @@ export * from '#src/drivers/DatabaseDriver'
export * from '#src/factories/ConnectionFactory'

export * from '#src/facades/Queue'
export * from '#src/worker/BaseWorker'
export * from '#src/worker/WorkerImpl'
export * from '#src/providers/QueueProvider'
export * from '#src/providers/WorkerProvider'
Expand Down
12 changes: 12 additions & 0 deletions src/types/ConnectionOptions.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -125,5 +125,17 @@ export type ConnectionOptions = {
* @default Parser.timeToMs('5m')
*/
workerTimeoutMs?: number

/**
* Define how many independent consumer loops run in parallel for this
* connection. Each loop still pulls and processes ONE job at a time, so a
* value of `N` yields an effective concurrency of `N`. This is honored by
* the `@athenna/event` consumer; raise it to drain a backed-up queue faster
* without spinning up extra processes. When the option is `null`/unset it
* defaults to `0`, which falls back to a single serial loop (one-by-one).
*
* @default Config.get(`queue.connections.${connection}.workerConcurrency`, 0)
*/
workerConcurrency?: number
}
}
7 changes: 5 additions & 2 deletions src/types/WorkerOptions.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -16,9 +16,12 @@ export type WorkerOptions = {
name?: string

/**
* Define how much instances of the same worker could run in parallel.
* Define how many instances of the same worker run in parallel. Each
* instance still processes one job at a time, so a value of `N` yields an
* effective concurrency of `N`. When omitted, the worker falls back to the
* connection's `workerConcurrency` config and, if that is `0`/unset, to `1`.
*
* @default 1
* @default Config.get(`queue.connections.${connection}.workerConcurrency`, 1)
*/
concurrency?: number

Expand Down
65 changes: 65 additions & 0 deletions src/worker/BaseWorker.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,65 @@
/**
* @athenna/queue
*
* (c) João Lenon <lenon@athenna.io>
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/

import 'reflect-metadata'

import { Queue } from '#src/facades/Queue'
import { Annotation } from '@athenna/ioc'
import type { QueueImpl } from '#src/queue/QueueImpl'

/**
* Base class for workers. Extend it to get a queue instance already bound to
* the worker's own connection, so you don't need to call `Queue.connection()`
* on every operation.
*
* @example
* ```ts
* @Worker()
* export class HelloWorker extends BaseWorker {
* public async handle(ctx: Context) {
* await this.queue.add({ hello: 'world' })
* }
* }
* ```
*/
export class BaseWorker {
/**
* Cached queue instance bound to this worker's connection.
*/
private _queue?: QueueImpl

/**
* The queue connection name of this worker. It is resolved from the
* worker's `@Worker({ connection })` metadata, falling back to the default
* connection (`queue.default`) when the worker is not annotated.
*/
public get connection() {
const meta = Annotation.getMeta(this.constructor)

return meta?.connection ?? Config.get('queue.default')
}

/**
* A queue instance already bound to this worker's connection. Use it to
* enqueue or inspect jobs without calling `Queue.connection(...)` on every
* operation.
*
* @example
* ```ts
* await this.queue.add({ email: 'lenon@athenna.io' })
* ```
*/
public get queue() {
if (!this._queue) {
this._queue = Queue.connection(this.connection)
}

return this._queue
}
}
25 changes: 24 additions & 1 deletion src/worker/WorkerTaskBuilder.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -200,13 +200,36 @@ export class WorkerTaskBuilder {
return
}

const n = this.worker.concurrency ?? 1
const n = this.resolveConcurrency()

for (let i = 0; i < n; i++) {
this.spawn()
}
}

/**
* Resolve how many worker loops to spawn. An explicit `concurrency` set on
* the worker (e.g. via `@Worker({ concurrency })`) wins, then the worker
* `options.workerConcurrency`, then the connection's `workerConcurrency`
* config. When none is a positive number it falls back to a single serial
* loop.
*/
private resolveConcurrency(): number {
const explicit =
this.worker.concurrency ?? this.worker.options?.workerConcurrency

if (Is.Number(explicit) && explicit > 0) {
return explicit
}

const configured = Config.get(
`queue.connections.${this.worker.connection}.workerConcurrency`,
0
)

return configured > 0 ? configured : 1
}

/**
* Use spawn to force a worker instance to run.
*/
Expand Down
4 changes: 2 additions & 2 deletions templates/worker.edge
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
import { Worker, type Context } from '@athenna/queue'
import { Worker, BaseWorker, type Context } from '@athenna/queue'

@Worker()
export class {{ namePascal }} {
export class {{ namePascal }} extends BaseWorker {
public async handle(ctx: Context) {
//
}
Expand Down
7 changes: 7 additions & 0 deletions tests/fixtures/config/queue.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -61,6 +61,13 @@ export default {
workerTimeoutMs: 200
},

memoryConcurrent: {
driver: 'memory',
queue: 'default',
deadletter: 'deadletter',
workerConcurrency: 3
},

aws_sqs: {
driver: 'aws_sqs',
type: 'standard',
Expand Down
82 changes: 82 additions & 0 deletions tests/unit/worker/BaseWorkerTest.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,82 @@
/**
* @athenna/queue
*
* (c) João Lenon <lenon@athenna.io>
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/

import { Path } from '@athenna/common'
import { Worker, BaseWorker, QueueImpl } from '#src'
import { LoggerProvider } from '@athenna/logger'
import { QueueProvider } from '#src/providers/QueueProvider'
import { Test, BeforeEach, AfterEach, type Context } from '@athenna/test'

@Worker({ connection: 'fake' })
class FakeConnectionWorker extends BaseWorker {
public async handle() {}
}

@Worker({ connection: 'memory' })
class MemoryConnectionWorker extends BaseWorker {
public async handle() {}
}

@Worker()
class DefaultConnectionWorker extends BaseWorker {
public async handle() {}
}

export class BaseWorkerTest {
@BeforeEach()
public async beforeEach() {
await Config.loadAll(Path.fixtures('config'))

new LoggerProvider().register()
new QueueProvider().register()
}

@AfterEach()
public async afterEach() {
await new QueueProvider().shutdown()

ioc.reconstruct()
Config.clear()
}

@Test()
public async shouldResolveTheConnectionFromTheWorkerMetadata({ assert }: Context) {
const worker = new FakeConnectionWorker()

assert.equal(worker.connection, 'fake')
assert.equal(worker.queue.connectionName, 'fake')
}

@Test()
public async shouldFallBackToTheDefaultConnectionWhenNotAnnotatedWithOne({ assert }: Context) {
const worker = new DefaultConnectionWorker()

assert.equal(worker.connection, Config.get('queue.default'))
assert.equal(worker.queue.connectionName, Config.get('queue.default'))
}

@Test()
public async shouldExposeAReadyToUseQueueInstanceBoundToTheWorkerConnection({ assert }: Context) {
const worker = new MemoryConnectionWorker()

assert.instanceOf(worker.queue, QueueImpl)
assert.equal(worker.queue.connectionName, 'memory')
assert.isTrue(worker.queue.isConnected())
}

@Test()
public async shouldCacheTheQueueInstanceBetweenAccesses({ assert }: Context) {
const worker = new FakeConnectionWorker()

const first = worker.queue
const second = worker.queue

assert.isTrue(first === second)
}
}
62 changes: 62 additions & 0 deletions tests/unit/worker/WorkerImplTest.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -176,6 +176,68 @@ export class WorkerImplTest {
assert.deepEqual(task?.worker.connection, 'fake')
}

@Test()
public async shouldSpawnASingleLoopWhenNoConcurrencyIsConfigured({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('default_concurrency')
.connection('memory')
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 1)

builder.stop()
}

@Test()
public async shouldSpawnConcurrentLoopsBasedOnConnectionWorkerConcurrencyConfig({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('config_concurrency')
.connection('memoryConcurrent')
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 3)

builder.stop()
}

@Test()
public async shouldSpawnConcurrentLoopsFromWorkerOptionsWorkerConcurrency({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('options_concurrency')
.connection('memory')
.options({ workerConcurrency: 4 })
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 4)

builder.stop()
}

@Test()
public async shouldLetExplicitConcurrencyOverrideTheConnectionWorkerConcurrencyConfig({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('explicit_concurrency')
.connection('memoryConcurrent')
.concurrency(5)
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 5)

builder.stop()
}

@Test()
public async shouldBeAbleToCreateAWorkerTaskWithCustomOptions({ assert }: Context) {
Queue.worker()
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
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
4 changes: 2 additions & 2 deletions package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion package.json
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
{
"name": "@athenna/queue",
"version": "5.32.0",
"version": "5.33.0",
"description": "The Athenna queue handler.",
"license": "MIT",
"author": "João Lenon <lenon@athenna.io>",
Expand Down
1 change: 1 addition & 0 deletions src/index.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -19,6 +19,7 @@ export * from '#src/drivers/DatabaseDriver'
export * from '#src/factories/ConnectionFactory'

export * from '#src/facades/Queue'
export * from '#src/worker/BaseWorker'
export * from '#src/worker/WorkerImpl'
export * from '#src/providers/QueueProvider'
export * from '#src/providers/WorkerProvider'
Expand Down
12 changes: 12 additions & 0 deletions src/types/ConnectionOptions.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -125,5 +125,17 @@ export type ConnectionOptions = {
* @default Parser.timeToMs('5m')
*/
workerTimeoutMs?: number

/**
* Define how many independent consumer loops run in parallel for this
* connection. Each loop still pulls and processes ONE job at a time, so a
* value of `N` yields an effective concurrency of `N`. This is honored by
* the `@athenna/event` consumer; raise it to drain a backed-up queue faster
* without spinning up extra processes. When the option is `null`/unset it
* defaults to `0`, which falls back to a single serial loop (one-by-one).
*
* @default Config.get(`queue.connections.${connection}.workerConcurrency`, 0)
*/
workerConcurrency?: number
}
}
7 changes: 5 additions & 2 deletions src/types/WorkerOptions.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -16,9 +16,12 @@ export type WorkerOptions = {
name?: string

/**
* Define how much instances of the same worker could run in parallel.
* Define how many instances of the same worker run in parallel. Each
* instance still processes one job at a time, so a value of `N` yields an
* effective concurrency of `N`. When omitted, the worker falls back to the
* connection's `workerConcurrency` config and, if that is `0`/unset, to `1`.
*
* @default 1
* @default Config.get(`queue.connections.${connection}.workerConcurrency`, 1)
*/
concurrency?: number

Expand Down
65 changes: 65 additions & 0 deletions src/worker/BaseWorker.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,65 @@
/**
* @athenna/queue
*
* (c) João Lenon <lenon@athenna.io>
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/

import 'reflect-metadata'

import { Queue } from '#src/facades/Queue'
import { Annotation } from '@athenna/ioc'
import type { QueueImpl } from '#src/queue/QueueImpl'

/**
* Base class for workers. Extend it to get a queue instance already bound to
* the worker's own connection, so you don't need to call `Queue.connection()`
* on every operation.
*
* @example
* ```ts
* @Worker()
* export class HelloWorker extends BaseWorker {
* public async handle(ctx: Context) {
* await this.queue.add({ hello: 'world' })
* }
* }
* ```
*/
export class BaseWorker {
/**
* Cached queue instance bound to this worker's connection.
*/
private _queue?: QueueImpl

/**
* The queue connection name of this worker. It is resolved from the
* worker's `@Worker({ connection })` metadata, falling back to the default
* connection (`queue.default`) when the worker is not annotated.
*/
public get connection() {
const meta = Annotation.getMeta(this.constructor)

return meta?.connection ?? Config.get('queue.default')
}

/**
* A queue instance already bound to this worker's connection. Use it to
* enqueue or inspect jobs without calling `Queue.connection(...)` on every
* operation.
*
* @example
* ```ts
* await this.queue.add({ email: 'lenon@athenna.io' })
* ```
*/
public get queue() {
if (!this._queue) {
this._queue = Queue.connection(this.connection)
}

return this._queue
}
}
25 changes: 24 additions & 1 deletion src/worker/WorkerTaskBuilder.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -200,13 +200,36 @@ export class WorkerTaskBuilder {
return
}

const n = this.worker.concurrency ?? 1
const n = this.resolveConcurrency()

for (let i = 0; i < n; i++) {
this.spawn()
}
}

/**
* Resolve how many worker loops to spawn. An explicit `concurrency` set on
* the worker (e.g. via `@Worker({ concurrency })`) wins, then the worker
* `options.workerConcurrency`, then the connection's `workerConcurrency`
* config. When none is a positive number it falls back to a single serial
* loop.
*/
private resolveConcurrency(): number {
const explicit =
this.worker.concurrency ?? this.worker.options?.workerConcurrency

if (Is.Number(explicit) && explicit > 0) {
return explicit
}

const configured = Config.get(
`queue.connections.${this.worker.connection}.workerConcurrency`,
0
)

return configured > 0 ? configured : 1
}

/**
* Use spawn to force a worker instance to run.
*/
Expand Down
4 changes: 2 additions & 2 deletions templates/worker.edge
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
import { Worker, type Context } from '@athenna/queue'
import { Worker, BaseWorker, type Context } from '@athenna/queue'

@Worker()
export class {{ namePascal }} {
export class {{ namePascal }} extends BaseWorker {
public async handle(ctx: Context) {
//
}
Expand Down
7 changes: 7 additions & 0 deletions tests/fixtures/config/queue.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -61,6 +61,13 @@ export default {
workerTimeoutMs: 200
},

memoryConcurrent: {
driver: 'memory',
queue: 'default',
deadletter: 'deadletter',
workerConcurrency: 3
},

aws_sqs: {
driver: 'aws_sqs',
type: 'standard',
Expand Down
82 changes: 82 additions & 0 deletions tests/unit/worker/BaseWorkerTest.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,82 @@
/**
* @athenna/queue
*
* (c) João Lenon <lenon@athenna.io>
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/

import { Path } from '@athenna/common'
import { Worker, BaseWorker, QueueImpl } from '#src'
import { LoggerProvider } from '@athenna/logger'
import { QueueProvider } from '#src/providers/QueueProvider'
import { Test, BeforeEach, AfterEach, type Context } from '@athenna/test'

@Worker({ connection: 'fake' })
class FakeConnectionWorker extends BaseWorker {
public async handle() {}
}

@Worker({ connection: 'memory' })
class MemoryConnectionWorker extends BaseWorker {
public async handle() {}
}

@Worker()
class DefaultConnectionWorker extends BaseWorker {
public async handle() {}
}

export class BaseWorkerTest {
@BeforeEach()
public async beforeEach() {
await Config.loadAll(Path.fixtures('config'))

new LoggerProvider().register()
new QueueProvider().register()
}

@AfterEach()
public async afterEach() {
await new QueueProvider().shutdown()

ioc.reconstruct()
Config.clear()
}

@Test()
public async shouldResolveTheConnectionFromTheWorkerMetadata({ assert }: Context) {
const worker = new FakeConnectionWorker()

assert.equal(worker.connection, 'fake')
assert.equal(worker.queue.connectionName, 'fake')
}

@Test()
public async shouldFallBackToTheDefaultConnectionWhenNotAnnotatedWithOne({ assert }: Context) {
const worker = new DefaultConnectionWorker()

assert.equal(worker.connection, Config.get('queue.default'))
assert.equal(worker.queue.connectionName, Config.get('queue.default'))
}

@Test()
public async shouldExposeAReadyToUseQueueInstanceBoundToTheWorkerConnection({ assert }: Context) {
const worker = new MemoryConnectionWorker()

assert.instanceOf(worker.queue, QueueImpl)
assert.equal(worker.queue.connectionName, 'memory')
assert.isTrue(worker.queue.isConnected())
}

@Test()
public async shouldCacheTheQueueInstanceBetweenAccesses({ assert }: Context) {
const worker = new FakeConnectionWorker()

const first = worker.queue
const second = worker.queue

assert.isTrue(first === second)
}
}
62 changes: 62 additions & 0 deletions tests/unit/worker/WorkerImplTest.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -176,6 +176,68 @@ export class WorkerImplTest {
assert.deepEqual(task?.worker.connection, 'fake')
}

@Test()
public async shouldSpawnASingleLoopWhenNoConcurrencyIsConfigured({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('default_concurrency')
.connection('memory')
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 1)

builder.stop()
}

@Test()
public async shouldSpawnConcurrentLoopsBasedOnConnectionWorkerConcurrencyConfig({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('config_concurrency')
.connection('memoryConcurrent')
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 3)

builder.stop()
}

@Test()
public async shouldSpawnConcurrentLoopsFromWorkerOptionsWorkerConcurrency({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('options_concurrency')
.connection('memory')
.options({ workerConcurrency: 4 })
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 4)

builder.stop()
}

@Test()
public async shouldLetExplicitConcurrencyOverrideTheConnectionWorkerConcurrencyConfig({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('explicit_concurrency')
.connection('memoryConcurrent')
.concurrency(5)
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 5)

builder.stop()
}

@Test()
public async shouldBeAbleToCreateAWorkerTaskWithCustomOptions({ assert }: Context) {
Queue.worker()
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content
Merged
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
4 changes: 2 additions & 2 deletions package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion package.json
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
{
"name": "@athenna/queue",
"version": "5.32.0",
"version": "5.33.0",
"description": "The Athenna queue handler.",
"license": "MIT",
"author": "João Lenon <lenon@athenna.io>",
Expand Down
1 change: 1 addition & 0 deletions src/index.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -19,6 +19,7 @@ export * from '#src/drivers/DatabaseDriver'
export * from '#src/factories/ConnectionFactory'

export * from '#src/facades/Queue'
export * from '#src/worker/BaseWorker'
export * from '#src/worker/WorkerImpl'
export * from '#src/providers/QueueProvider'
export * from '#src/providers/WorkerProvider'
Expand Down
12 changes: 12 additions & 0 deletions src/types/ConnectionOptions.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -125,5 +125,17 @@ export type ConnectionOptions = {
* @default Parser.timeToMs('5m')
*/
workerTimeoutMs?: number

/**
* Define how many independent consumer loops run in parallel for this
* connection. Each loop still pulls and processes ONE job at a time, so a
* value of `N` yields an effective concurrency of `N`. This is honored by
* the `@athenna/event` consumer; raise it to drain a backed-up queue faster
* without spinning up extra processes. When the option is `null`/unset it
* defaults to `0`, which falls back to a single serial loop (one-by-one).
*
* @default Config.get(`queue.connections.${connection}.workerConcurrency`, 0)
*/
workerConcurrency?: number
}
}
7 changes: 5 additions & 2 deletions src/types/WorkerOptions.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -16,9 +16,12 @@ export type WorkerOptions = {
name?: string

/**
* Define how much instances of the same worker could run in parallel.
* Define how many instances of the same worker run in parallel. Each
* instance still processes one job at a time, so a value of `N` yields an
* effective concurrency of `N`. When omitted, the worker falls back to the
* connection's `workerConcurrency` config and, if that is `0`/unset, to `1`.
*
* @default 1
* @default Config.get(`queue.connections.${connection}.workerConcurrency`, 1)
*/
concurrency?: number

Expand Down
65 changes: 65 additions & 0 deletions src/worker/BaseWorker.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,65 @@
/**
* @athenna/queue
*
* (c) João Lenon <lenon@athenna.io>
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/

import 'reflect-metadata'

import { Queue } from '#src/facades/Queue'
import { Annotation } from '@athenna/ioc'
import type { QueueImpl } from '#src/queue/QueueImpl'

/**
* Base class for workers. Extend it to get a queue instance already bound to
* the worker's own connection, so you don't need to call `Queue.connection()`
* on every operation.
*
* @example
* ```ts
* @Worker()
* export class HelloWorker extends BaseWorker {
* public async handle(ctx: Context) {
* await this.queue.add({ hello: 'world' })
* }
* }
* ```
*/
export class BaseWorker {
/**
* Cached queue instance bound to this worker's connection.
*/
private _queue?: QueueImpl

/**
* The queue connection name of this worker. It is resolved from the
* worker's `@Worker({ connection })` metadata, falling back to the default
* connection (`queue.default`) when the worker is not annotated.
*/
public get connection() {
const meta = Annotation.getMeta(this.constructor)

return meta?.connection ?? Config.get('queue.default')
}

/**
* A queue instance already bound to this worker's connection. Use it to
* enqueue or inspect jobs without calling `Queue.connection(...)` on every
* operation.
*
* @example
* ```ts
* await this.queue.add({ email: 'lenon@athenna.io' })
* ```
*/
public get queue() {
if (!this._queue) {
this._queue = Queue.connection(this.connection)
}

return this._queue
}
}
25 changes: 24 additions & 1 deletion src/worker/WorkerTaskBuilder.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -200,13 +200,36 @@ export class WorkerTaskBuilder {
return
}

const n = this.worker.concurrency ?? 1
const n = this.resolveConcurrency()

for (let i = 0; i < n; i++) {
this.spawn()
}
}

/**
* Resolve how many worker loops to spawn. An explicit `concurrency` set on
* the worker (e.g. via `@Worker({ concurrency })`) wins, then the worker
* `options.workerConcurrency`, then the connection's `workerConcurrency`
* config. When none is a positive number it falls back to a single serial
* loop.
*/
private resolveConcurrency(): number {
const explicit =
this.worker.concurrency ?? this.worker.options?.workerConcurrency

if (Is.Number(explicit) && explicit > 0) {
return explicit
}

const configured = Config.get(
`queue.connections.${this.worker.connection}.workerConcurrency`,
0
)

return configured > 0 ? configured : 1
}

/**
* Use spawn to force a worker instance to run.
*/
Expand Down
4 changes: 2 additions & 2 deletions templates/worker.edge
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
import { Worker, type Context } from '@athenna/queue'
import { Worker, BaseWorker, type Context } from '@athenna/queue'

@Worker()
export class {{ namePascal }} {
export class {{ namePascal }} extends BaseWorker {
public async handle(ctx: Context) {
//
}
Expand Down
7 changes: 7 additions & 0 deletions tests/fixtures/config/queue.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -61,6 +61,13 @@ export default {
workerTimeoutMs: 200
},

memoryConcurrent: {
driver: 'memory',
queue: 'default',
deadletter: 'deadletter',
workerConcurrency: 3
},

aws_sqs: {
driver: 'aws_sqs',
type: 'standard',
Expand Down
82 changes: 82 additions & 0 deletions tests/unit/worker/BaseWorkerTest.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,82 @@
/**
* @athenna/queue
*
* (c) João Lenon <lenon@athenna.io>
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/

import { Path } from '@athenna/common'
import { Worker, BaseWorker, QueueImpl } from '#src'
import { LoggerProvider } from '@athenna/logger'
import { QueueProvider } from '#src/providers/QueueProvider'
import { Test, BeforeEach, AfterEach, type Context } from '@athenna/test'

@Worker({ connection: 'fake' })
class FakeConnectionWorker extends BaseWorker {
public async handle() {}
}

@Worker({ connection: 'memory' })
class MemoryConnectionWorker extends BaseWorker {
public async handle() {}
}

@Worker()
class DefaultConnectionWorker extends BaseWorker {
public async handle() {}
}

export class BaseWorkerTest {
@BeforeEach()
public async beforeEach() {
await Config.loadAll(Path.fixtures('config'))

new LoggerProvider().register()
new QueueProvider().register()
}

@AfterEach()
public async afterEach() {
await new QueueProvider().shutdown()

ioc.reconstruct()
Config.clear()
}

@Test()
public async shouldResolveTheConnectionFromTheWorkerMetadata({ assert }: Context) {
const worker = new FakeConnectionWorker()

assert.equal(worker.connection, 'fake')
assert.equal(worker.queue.connectionName, 'fake')
}

@Test()
public async shouldFallBackToTheDefaultConnectionWhenNotAnnotatedWithOne({ assert }: Context) {
const worker = new DefaultConnectionWorker()

assert.equal(worker.connection, Config.get('queue.default'))
assert.equal(worker.queue.connectionName, Config.get('queue.default'))
}

@Test()
public async shouldExposeAReadyToUseQueueInstanceBoundToTheWorkerConnection({ assert }: Context) {
const worker = new MemoryConnectionWorker()

assert.instanceOf(worker.queue, QueueImpl)
assert.equal(worker.queue.connectionName, 'memory')
assert.isTrue(worker.queue.isConnected())
}

@Test()
public async shouldCacheTheQueueInstanceBetweenAccesses({ assert }: Context) {
const worker = new FakeConnectionWorker()

const first = worker.queue
const second = worker.queue

assert.isTrue(first === second)
}
}
62 changes: 62 additions & 0 deletions tests/unit/worker/WorkerImplTest.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -176,6 +176,68 @@ export class WorkerImplTest {
assert.deepEqual(task?.worker.connection, 'fake')
}

@Test()
public async shouldSpawnASingleLoopWhenNoConcurrencyIsConfigured({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('default_concurrency')
.connection('memory')
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 1)

builder.stop()
}

@Test()
public async shouldSpawnConcurrentLoopsBasedOnConnectionWorkerConcurrencyConfig({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('config_concurrency')
.connection('memoryConcurrent')
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 3)

builder.stop()
}

@Test()
public async shouldSpawnConcurrentLoopsFromWorkerOptionsWorkerConcurrency({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('options_concurrency')
.connection('memory')
.options({ workerConcurrency: 4 })
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 4)

builder.stop()
}

@Test()
public async shouldLetExplicitConcurrencyOverrideTheConnectionWorkerConcurrencyConfig({ assert }: Context) {
const builder = Queue.worker()
.task()
.name('explicit_concurrency')
.connection('memoryConcurrent')
.concurrency(5)
.handler(() => {})

builder.start()

assert.lengthOf((builder as any).timers, 5)

builder.stop()
}

@Test()
public async shouldBeAbleToCreateAWorkerTaskWithCustomOptions({ assert }: Context) {
Queue.worker()
Expand Down
Loading