diff --git a/README.md b/README.md index 4f166a6..79218c1 100644 --- a/README.md +++ b/README.md @@ -48,13 +48,21 @@ paths concurrently. clean up items written to a queue by a process outside your control (e.g., webhooks). * `healthyPingLatency: {number | string}` the maximum response latency to pings that is considered "healthy" for this queue. + * `captureLeaseTransactionMetrics: {function(string, number, number)}` a callback invoked after + each acquired, contended, or failed task lease transaction. It receives the acquisition + outcome, NodeFire transaction tries, and transaction duration in milliseconds. Missing optional + NodeFire metadata is reported as zero. The callback must be synchronous. Callback errors are + reported through + `settings.captureError` and do not affect task processing. * `@param {function(Object):RETRY | number | string | undefined}` worker The worker function that handles enqueued tasks. It will be given a task object as argument, with a special $ref attribute - set to the Nodefire ref of that task. The worker can perform arbitrary computation whose duration - should not exceed the queue's minLease value. It can manipulate the task itself in Firebase as - well, e.g. to delete it (to get at-most-once queue semantics) or otherwise modify it. The worker - can return any of the following: + set to the Nodefire ref of that task. On a task's first acquisition, the worker-facing `_lease` + object also has a non-enumerable `firstAcquisition: true` property that is not saved to Firebase; + the property is absent on subsequent acquisitions. The worker can perform arbitrary computation + whose duration should not exceed the queue's minLease value. It can manipulate the task itself in + Firebase as well, e.g. to delete it (to get at-most-once queue semantics) or otherwise modify it. + The worker can return any of the following: * undefined or null to cause the task to be retired from the queue. * firelease.RETRY to cause the task to be retried after the current lease expires (and reset the lease backoff counter). @@ -142,8 +150,16 @@ Mutable default option values for all subsequent attachWorker calls. See that f ```stats: {Object}``` The live stats object also passed to the ping callback. Global fields include `healthy`, -`sickQueues`, `sickSources`, `stuckTasks`, `maxLatency`, and `tasksAcquired`. Each entry in `queues` -includes its own health, latency, acquisition count, and all physical `sources`. Queue-level +`sickQueues`, `sickSources`, `stuckTasks`, `maxLatency`, `leaseTransactions`, and the legacy +`tasksAcquired`. `leaseTransactions` contains lifetime `acquired`, `contended`, `failed`, and +`tries` counts for task lease transactions. `tries` includes failed transactions and comes from +NodeFire transaction metadata. `duration` is an exponential moving average of the NodeFire +transaction duration in milliseconds, using an alpha of 0.1. These stats are available for every +physical source. Logical queue and global counts are additive, while their duration is the average +of the underlying source or queue duration values that have recorded at least one attempt. The +legacy `tasksAcquired` field also remains +lifetime-cumulative. Each entry in `queues` includes its own health, latency, leasing totals, and +all physical `sources`. Queue-level `size` and `sizeDelta` sum their source values when all are known, `sizeTimestamp` is the oldest source timestamp, and `mode` is `full`, `safe`, or `mixed`. Source stats include `connected`, current `mode` (`full` or `safe`), last known `size`, `sizeTimestamp` when the size came from a diff --git a/package.json b/package.json index b249831..94c0119 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "firelease", - "version": "4.1.0", + "version": "4.2.0", "packageManager": "yarn@4.13.0", "description": "Firebase queue consumer for Node with at-least-once semantics", "main": "built/index.js", diff --git a/src/firelease.ts b/src/firelease.ts index 81f720d..4e763f3 100644 --- a/src/firelease.ts +++ b/src/firelease.ts @@ -1,6 +1,6 @@ import _ from 'lodash'; import ms from 'ms'; -import NodeFire from 'nodefire'; +import NodeFire, {type TransactionMetadata} from 'nodefire'; import * as timers from 'safe-timers'; import { FireleaseStats, QueueSourceStats, QueueStats, type QueueSourceMode @@ -14,6 +14,7 @@ const QUEUE_CHECK_TIMEOUT = ms('15s'); const QUEUE_SIZE_HYSTERESIS = 0.15; const QUEUE_SIZE_MISMATCH_THRESHOLD = 100; const DEMOTION_JITTER = ms('30s'); +const LEASE_TRANSACTION_DURATION_ALPHA = 0.1; declare const RETRY_DIRECTIVE: unique symbol; @@ -29,6 +30,10 @@ export interface Lease { extendLeasePromise?: Promise; } +export type AcquiredLease = Lease & { + expiry: number, time: number, attempts: number, initial: number, readonly firstAcquisition?: true +}; + export interface LeaseItem { _lease?: Lease; [key: string]: any; @@ -39,7 +44,7 @@ export interface RetryDirective { } export interface WorkerItem extends LeaseItem { - _lease: Lease & {expiry: number}; + _lease: AcquiredLease; readonly $ref: NodeFire; readonly $leaseTimeRemaining: number; } @@ -69,6 +74,11 @@ export interface FireleaseError extends Error { level?: FireleaseErrorLevel; } +export type LeaseTransactionOutcome = 'acquired' | 'contended' | 'failed'; +export type CaptureLeaseTransactionMetrics = ( + outcome: LeaseTransactionOutcome, tries: number, duration: number +) => undefined; + export interface QueueOptions { maxConcurrent?: number; bufferSize?: number; @@ -76,6 +86,7 @@ export interface QueueOptions { maxLease?: Duration; healthyPingLatency?: Duration; preprocess?: (item: LeaseItem) => LeaseItem; + captureLeaseTransactionMetrics?: CaptureLeaseTransactionMetrics; } export type PingReport = FireleaseStats; @@ -121,6 +132,7 @@ interface NormalizedQueueOptions { maxLease: number; healthyPingLatency: number; preprocess?: (item: LeaseItem) => LeaseItem; + captureLeaseTransactionMetrics?: CaptureLeaseTransactionMetrics; } const queues: Queue[] = []; @@ -261,12 +273,16 @@ class Task { async process() { let startTimestamp = 0; let acquired = false; + let contended = false; let reschedule = true; + let firstAcquisition = false; this.working = true; this.phase = 'lease'; const transactionPromise = this.ref.transaction(itemValue => { const item = itemValue as LeaseItem | null; acquired = false; + contended = false; + firstAcquisition = false; if (tasks[this.key] !== this || this.removed) return; if (!item || this.ref.key === PING_KEY) { acquired = true; @@ -276,9 +292,11 @@ class Task { // console.log('txn ', this.ref.key, 'lease', item._lease, 'now', startTimestamp); // Check if another process beat us to it. if (item._lease?.expiry && item._lease.expiry > startTimestamp) { + contended = true; return item; } acquired = true; + firstAcquisition = _.isNil(item._lease?.initial); item._lease ??= {}; item._lease.time = this.queue.constrainLeaseDuration((item._lease.time ?? 0) * 2); item._lease.expiry = startTimestamp + item._lease.time; @@ -287,14 +305,22 @@ class Task { item._lease.busy = true; return this.queue.callPreprocess(item); }, {detectStuck: 5, prefetchValue: false, timeout: ms('15s')}); + let transactionCompleted = false; try { const item = await transactionPromise; + transactionCompleted = true; if (acquired && item !== null && this.ref.key !== PING_KEY) { - if (!_.isObject(item)) throw new Error(`item not an object: ${item}`); + this.recordLeaseTransaction('acquired', transactionPromise.transaction); + if (firstAcquisition) Object.defineProperty(item._lease, 'firstAcquisition', {value: true}); this.queue.stats.tasksAcquired++; await this.run(item as WorkerItem, startTimestamp); + } else if (contended) { + this.recordLeaseTransaction('contended', transactionPromise.transaction); } } catch (error) { + if (!transactionCompleted && this.ref.key !== PING_KEY) { + this.recordLeaseTransaction('failed', transactionPromise.transaction); + } reschedule = false; // Hardcoded retry -- hard to do anything smarter, since we failed to update the task in // Firebase. @@ -315,6 +341,37 @@ class Task { } } + recordLeaseTransaction( + outcome: LeaseTransactionOutcome, transaction: TransactionMetadata | undefined + ) { + try { + const tries = transaction?.tries ?? 0; + const transactionDuration = transaction?.duration ?? 0; + const leaseStats = this.source.stats.leaseTransactions; + const priorAttempts = leaseStats.acquired + leaseStats.contended + leaseStats.failed; + leaseStats[outcome] += 1; + leaseStats.tries += tries; + leaseStats.duration = priorAttempts === 0 ? + transactionDuration : + leaseStats.duration * (1 - LEASE_TRANSACTION_DURATION_ALPHA) + + transactionDuration * LEASE_TRANSACTION_DURATION_ALPHA; + this.queue.options.captureLeaseTransactionMetrics?.(outcome, tries, transactionDuration); + } catch (error) { + try { + const metricError: FireleaseError = _.isError(error) ? error : new Error(String(error)); + metricError.firelease = _.assign( + metricError.firelease ?? {}, {itemKey: this.key, phase: 'lease-metric'}); + settings.captureError(metricError); + } catch (captureError) { + try { + console.error('Error capturing lease transaction metric error:', captureError); + } catch { + // Metric recording must never interrupt task processing. + } + } + } + } + async run(item: WorkerItem, startTimestamp: number) { Object.defineProperty(item, '$ref', {value: this.ref}); Object.defineProperty(item, '$leaseTimeRemaining', {get: () => { @@ -395,7 +452,7 @@ class Task { if (currentItem._lease) delete currentItem._lease.busy; return currentItem; }, {prefetchValue: false}) as LeaseItem | null | undefined; - if (item2) item._lease = item2._lease as Lease & {expiry: number}; + if (item2) item._lease = item2._lease as AcquiredLease; } catch (postProcessingError) { this.handlePostProcessingError(postProcessingError); } @@ -873,7 +930,7 @@ class QueueSource { } recordPingResult(startedAt: number, succeeded: boolean) { - const latency = performance.now() - startedAt; + const latency = Math.round(performance.now() - startedAt); this.stats.latency = latency; this.stats.healthy = succeeded && latency < this.queue.options.healthyPingLatency; this.stats.pingTimestamp = Date.now(); @@ -1069,9 +1126,15 @@ class Queue { * control (e.g., webhooks). * healthyPingLatency: {number | string} the maximum response latency to pings that is * considered "healthy" for this queue. + * captureLeaseTransactionMetrics: {function(string, number, number)} a callback invoked + * after each acquired, contended, or failed task lease transaction with its outcome, + * NodeFire transaction tries, and duration in milliseconds. The callback must be + * synchronous. * @param {function(Object):RETRY | number | string | undefined} worker The worker function that * handles enqueued tasks. It will be given a task object as argument, with a special $ref - * attribute set to the Nodefire ref of that task. The worker can perform arbitrary + * attribute set to the Nodefire ref of that task. On a task's first acquisition its _lease + * also has a non-enumerable firstAcquisition property set to true; it is not saved to + * Firebase and is absent on subsequent acquisitions. The worker can perform arbitrary * computation whose duration should not exceed the queue's minLease value. It can * manipulate the task itself in Firebase as well, e.g. to delete it (to get at-most-once * queue semantics) or otherwise modify it. The worker can return any of the following: diff --git a/src/stats.ts b/src/stats.ts index f949faa..95f91f9 100644 --- a/src/stats.ts +++ b/src/stats.ts @@ -3,6 +3,33 @@ import _ from 'lodash'; export type QueueSourceMode = 'full' | 'safe'; export type QueueMode = QueueSourceMode | 'mixed'; +export interface LeaseTransactionStats { + acquired: number; + contended: number; + failed: number; + tries: number; + duration: number; +} + +function createLeaseTransactionStats(): LeaseTransactionStats { + return {acquired: 0, contended: 0, failed: 0, tries: 0, duration: 0}; +} + +function rollUpLeaseTransactions(items: LeaseTransactionStats[]) { + const result = createLeaseTransactionStats(); + result.acquired = _.sumBy(items, 'acquired'); + result.contended = _.sumBy(items, 'contended'); + result.failed = _.sumBy(items, 'failed'); + result.tries = _.sumBy(items, 'tries'); + const attemptedItems = _.filter(items, countLeaseAttempts); + if (attemptedItems.length) result.duration = _.meanBy(attemptedItems, 'duration'); + return result; +} + +function countLeaseAttempts(stats: LeaseTransactionStats) { + return stats.acquired + stats.contended + stats.failed; +} + function exposeGetters(instance: object, properties: string[]) { const prototype = Object.getPrototypeOf(instance); for (const property of properties) { @@ -20,6 +47,7 @@ export class QueueSourceStats { healthy = true; latency: number | null = null; declare pingTimestamp?: number; + readonly leaseTransactions = createLeaseTransactionStats(); constructor(readonly ref: string) {} } @@ -32,7 +60,10 @@ export class QueueStats { readonly key: string | null, readonly sources: QueueSourceStats[] ) { - exposeGetters(this, ['mode', 'size', 'sizeDelta', 'sizeTimestamp', 'healthy', 'maxLatency']); + exposeGetters( + this, + ['mode', 'size', 'sizeDelta', 'sizeTimestamp', 'healthy', 'maxLatency', 'leaseTransactions'], + ); } get mode(): QueueMode { @@ -61,6 +92,10 @@ export class QueueStats { get maxLatency() { return _(this.sources).map('latency').max() || 0; } + + get leaseTransactions() { + return rollUpLeaseTransactions(_.map(this.sources, 'leaseTransactions')); + } } export class FireleaseStats { @@ -70,7 +105,11 @@ export class FireleaseStats { constructor(getStuckTasks: () => number) { this.#getStuckTasks = getStuckTasks; exposeGetters( - this, ['healthy', 'sickQueues', 'sickSources', 'stuckTasks', 'maxLatency', 'tasksAcquired'], + this, + [ + 'healthy', 'sickQueues', 'sickSources', 'stuckTasks', 'maxLatency', 'leaseTransactions', + 'tasksAcquired' + ], ); } @@ -99,6 +138,10 @@ export class FireleaseStats { return _(this.queues).map('maxLatency').max() || 0; } + get leaseTransactions() { + return rollUpLeaseTransactions(_.map(this.queues, 'leaseTransactions')); + } + get tasksAcquired() { return _.sumBy(this.queues, queue => queue.tasksAcquired); } diff --git a/tests/fake_firebase.ts b/tests/fake_firebase.ts index c44656b..435bef9 100644 --- a/tests/fake_firebase.ts +++ b/tests/fake_firebase.ts @@ -1,5 +1,6 @@ import assert from 'node:assert'; import type NodeFire from 'nodefire'; +import type {TransactionMetadata} from 'nodefire'; interface FakeLease { [key: string]: unknown; @@ -92,12 +93,24 @@ export class FakeTaskRef { transaction( update: (value: FakeTaskValue | null) => FakeTaskValue | null | undefined ) { - if (this.queueRef.transactionError) return Promise.reject(this.queueRef.transactionError); - const previous = clone(this.value); - const updated = update(clone(this.value)); - if (updated !== undefined) this.value = clone(updated); - this.queueRef.notifyTaskChange(this, previous); - return Promise.resolve(clone(this.value)); + this.queueRef.beforeTransaction?.(this); + const metadata: TransactionMetadata = { + outcome: this.queueRef.transactionError ? 'error' : 'commit', + tries: this.queueRef.transactionTries, + duration: this.queueRef.transactionDuration + }; + let transactionPromise: Promise; + if (this.queueRef.transactionError) { + transactionPromise = Promise.reject(this.queueRef.transactionError); + } else { + const previous = clone(this.value); + const updated = update(clone(this.value)); + if (updated !== undefined) this.value = clone(updated); + this.queueRef.notifyTaskChange(this, previous); + transactionPromise = Promise.resolve(clone(this.value)); + } + return this.queueRef.omitTransactionMetadata ? + transactionPromise : Object.assign(transactionPromise, {transaction: metadata}); } get() { @@ -214,6 +227,10 @@ export class FakeQueueRef { childrenKeysCalls = 0; listenerError?: Error; transactionError?: Error; + omitTransactionMetadata = false; + transactionTries = 1; + transactionDuration = 0; + beforeTransaction?: (ref: FakeTaskRef) => void; fixedNow?: number; constructor(databaseName: string, readonly path: string) { diff --git a/tests/first_acquisition.test.ts b/tests/first_acquisition.test.ts new file mode 100644 index 0000000..da9140c --- /dev/null +++ b/tests/first_acquisition.test.ts @@ -0,0 +1,61 @@ +import assert from 'node:assert'; +import {test} from 'node:test'; + +import firelease, {TESTABLES} from '../src'; +import {asNodeFire, FakeQueueRef, waitFor} from './fake_firebase'; + +test('firstAcquisition is worker-local and only set on the first lease', async () => { + TESTABLES.resetBetweenTests(); + try { + const source = new FakeQueueRef('first-acquisition-database', 'queues/jobs'); + const flags: (true | undefined)[] = []; + let firstDescriptor: PropertyDescriptor | undefined; + let releaseFirst!: () => void; + const firstBlocked = new Promise(resolve => {releaseFirst = resolve;}); + firelease.attachWorker( + asNodeFire(source), + {minLease: 100, maxLease: 100}, + async item => { + flags.push(item._lease.firstAcquisition); + if (flags.length === 1) { + firstDescriptor = Object.getOwnPropertyDescriptor(item._lease, 'firstAcquisition'); + await firstBlocked; + return firelease.RETRY; + } + } + ); + + const task = source.addTask('task', {payload: 1}); + await waitFor(() => flags.length === 1); + + assert.deepStrictEqual(firstDescriptor, { + value: true, writable: false, enumerable: false, configurable: false + }); + assert.strictEqual(task.value?._lease?.firstAcquisition, undefined); + + releaseFirst(); + await waitFor(() => flags.length === 2 && !task.value); + assert.deepStrictEqual(flags, [true, undefined]); + } finally { + TESTABLES.resetBetweenTests(); + } +}); + +test('an existing initial value is not treated as a first acquisition', async () => { + TESTABLES.resetBetweenTests(); + try { + const source = new FakeQueueRef('existing-initial-database', 'queues/jobs'); + let workerCalled = false; + let firstAcquisition: true | undefined; + firelease.attachWorker(asNodeFire(source), item => { + workerCalled = true; + firstAcquisition = item._lease.firstAcquisition; + }); + + const task = source.addTask('task', {payload: 1, _lease: {initial: 0, expiry: 0}}); + await waitFor(() => workerCalled && !task.value); + assert.strictEqual(firstAcquisition, undefined); + } finally { + TESTABLES.resetBetweenTests(); + } +}); diff --git a/tests/lease_stats.test.ts b/tests/lease_stats.test.ts new file mode 100644 index 0000000..579fa93 --- /dev/null +++ b/tests/lease_stats.test.ts @@ -0,0 +1,104 @@ +import assert from 'node:assert'; +import {test} from 'node:test'; + +import firelease, {TESTABLES, type LeaseTransactionOutcome} from '../src'; +import {asNodeFire, FakeQueueRef, waitFor} from './fake_firebase'; + +test('lease stats capture all transaction outcomes through the hierarchy', async () => { + TESTABLES.resetBetweenTests(); + try { + const source = new FakeQueueRef('lease-stats-database', 'queues/jobs'); + source.transactionTries = 3; + source.transactionDuration = 25; + const capturedMetrics: [LeaseTransactionOutcome, number, number][] = []; + const capturedErrors: Error[] = []; + firelease.settings.captureError = error => {capturedErrors.push(error);}; + let workerCalls = 0; + firelease.attachWorker( + asNodeFire(source), + { + captureLeaseTransactionMetrics: (outcome, tries, duration) => { + capturedMetrics.push([outcome, tries, duration]); + if (outcome === 'acquired') throw new Error('Metric capture failed'); + } + }, + item => { + workerCalls++; + assert.strictEqual(item.payload, 'acquired'); + } + ); + const queueStats = firelease.stats.queues[firelease.stats.queues.length - 1]; + const sourceStats = queueStats.sources[0]; + + const acquiredTask = source.addTask('acquired', {payload: 'acquired'}); + await waitFor(() => workerCalls === 1 && !acquiredTask.value); + + source.transactionTries = 2; + source.transactionDuration = 10; + source.beforeTransaction = task => { + source.beforeTransaction = undefined; + task.value!._lease = {expiry: source.now + 60_000}; + }; + source.addTask('contended', {payload: 'contended'}); + await waitFor(() => sourceStats.leaseTransactions.contended === 1); + + source.transactionTries = 4; + source.transactionDuration = 50; + source.transactionError = new Error('Lease transaction failed'); + source.addTask('failed', {payload: 'failed'}); + await waitFor(() => capturedMetrics.length === 3); + + const expectedCounts = {acquired: 1, contended: 1, failed: 1, tries: 9}; + assert.strictEqual(workerCalls, 1); + assert.deepStrictEqual(capturedMetrics, [ + ['acquired', 3, 25], + ['contended', 2, 10], + ['failed', 4, 50] + ]); + assert.deepStrictEqual( + capturedErrors.map(error => error.message), + ['Metric capture failed', 'Lease transaction failed'], + ); + for (const stats of [ + sourceStats.leaseTransactions, queueStats.leaseTransactions, firelease.stats.leaseTransactions + ]) { + const {duration, ...counts} = stats; + assert.deepStrictEqual(counts, expectedCounts); + assert.ok(Math.abs(duration - 26.15) < Number.EPSILON * 26.15); + } + assert.strictEqual(queueStats.tasksAcquired, 1); + assert.strictEqual(firelease.stats.tasksAcquired, 1); + } finally { + TESTABLES.resetBetweenTests(); + } +}); + +test('metric recording errors cannot interrupt task processing', async () => { + TESTABLES.resetBetweenTests(); + try { + const source = new FakeQueueRef('metric-error-database', 'queues/jobs'); + source.omitTransactionMetadata = true; + let captureErrorCalls = 0; + firelease.settings.captureError = () => { + captureErrorCalls++; + throw new Error('Error capture failed'); + }; + let workerCalls = 0; + firelease.attachWorker( + asNodeFire(source), + {captureLeaseTransactionMetrics: () => {throw new Error('Metric capture failed');}}, + () => {workerCalls++;}, + ); + const sourceStats = firelease.stats.queues[firelease.stats.queues.length - 1].sources[0]; + + const task = source.addTask('acquired', {payload: 'acquired'}); + await waitFor(() => workerCalls === 1 && !task.value); + + assert.strictEqual(captureErrorCalls, 1); + assert.deepStrictEqual(sourceStats.leaseTransactions, { + acquired: 1, contended: 0, failed: 0, tries: 0, duration: 0 + }); + } finally { + TESTABLES.resetBetweenTests(); + } +}); diff --git a/tests/stats.test.ts b/tests/stats.test.ts index 2192754..f136e1f 100644 --- a/tests/stats.test.ts +++ b/tests/stats.test.ts @@ -5,6 +5,7 @@ import {FireleaseStats, QueueSourceStats, QueueStats, TESTABLES} from '../src'; test('stats are derived through the hierarchy on demand', () => { TESTABLES.resetBetweenTests(); + const expectedDuration = 40; let stuckTasks = 0; const sourceStats = new QueueSourceStats('https://stats.example.test/queues/jobs'); const secondSourceStats = new QueueSourceStats('https://second.example.test/queues/jobs'); @@ -27,11 +28,21 @@ test('stats are derived through the hierarchy on demand', () => { sourceStats.sizeDelta = 1; sourceStats.sizeTimestamp = 200; sourceStats.latency = 12; + sourceStats.leaseTransactions.acquired = 2; + sourceStats.leaseTransactions.contended = 3; + sourceStats.leaseTransactions.failed = 1; + sourceStats.leaseTransactions.tries = 8; + sourceStats.leaseTransactions.duration = 50; secondSourceStats.connected = true; secondSourceStats.size = 6; secondSourceStats.sizeDelta = 2; secondSourceStats.sizeTimestamp = 100; secondSourceStats.latency = 8; + secondSourceStats.leaseTransactions.acquired = 1; + secondSourceStats.leaseTransactions.contended = 2; + secondSourceStats.leaseTransactions.failed = 2; + secondSourceStats.leaseTransactions.tries = 5; + secondSourceStats.leaseTransactions.duration = 30; queueStats.tasksAcquired = 3; stuckTasks = 2; assert.strictEqual(hierarchyStats.healthy, true); @@ -39,6 +50,10 @@ test('stats are derived through the hierarchy on demand', () => { assert.deepStrictEqual(hierarchyStats.sickSources, []); assert.strictEqual(queueStats.maxLatency, 12); assert.strictEqual(hierarchyStats.maxLatency, 12); + assert.deepStrictEqual(queueStats.leaseTransactions, { + acquired: 3, contended: 5, failed: 3, tries: 13, duration: expectedDuration + }); + assert.deepStrictEqual(hierarchyStats.leaseTransactions, queueStats.leaseTransactions); assert.strictEqual(hierarchyStats.tasksAcquired, 3); assert.strictEqual(hierarchyStats.stuckTasks, 2); assert.strictEqual(queueStats.size, 10); @@ -52,7 +67,10 @@ test('stats are derived through the hierarchy on demand', () => { secondSourceStats.sizeDelta = 2; assert.deepStrictEqual( Object.keys(hierarchyStats), - ['queues', 'healthy', 'sickQueues', 'sickSources', 'stuckTasks', 'maxLatency', 'tasksAcquired'], + [ + 'queues', 'healthy', 'sickQueues', 'sickSources', 'stuckTasks', 'maxLatency', + 'leaseTransactions', 'tasksAcquired' + ], ); assert.deepStrictEqual(JSON.parse(JSON.stringify(hierarchyStats)), { queues: [{ @@ -65,6 +83,7 @@ test('stats are derived through the hierarchy on demand', () => { size: 4, healthy: true, latency: 12, + leaseTransactions: {acquired: 2, contended: 3, failed: 1, tries: 8, duration: 50}, ref: 'https://stats.example.test/queues/jobs', sizeDelta: 1, sizeTimestamp: 200 @@ -74,6 +93,7 @@ test('stats are derived through the hierarchy on demand', () => { size: 6, healthy: true, latency: 8, + leaseTransactions: {acquired: 1, contended: 2, failed: 2, tries: 5, duration: 30}, ref: 'https://second.example.test/queues/jobs', sizeTimestamp: 100, sizeDelta: 2 @@ -83,14 +103,49 @@ test('stats are derived through the hierarchy on demand', () => { sizeDelta: 3, sizeTimestamp: 100, healthy: true, - maxLatency: 12 + maxLatency: 12, + leaseTransactions: { + acquired: 3, contended: 5, failed: 3, tries: 13, duration: expectedDuration + } }], healthy: true, sickQueues: [], sickSources: [], stuckTasks: 2, maxLatency: 12, + leaseTransactions: { + acquired: 3, contended: 5, failed: 3, tries: 13, duration: expectedDuration + }, tasksAcquired: 3 }); TESTABLES.resetBetweenTests(); }); + +test('duration rollups use an unweighted mean at each level', () => { + const firstSource = new QueueSourceStats('https://first.example.test/queues/jobs'); + firstSource.leaseTransactions.acquired = 1; + firstSource.leaseTransactions.contended = 2; + firstSource.leaseTransactions.duration = 20; + const secondSource = new QueueSourceStats('https://second.example.test/queues/jobs'); + secondSource.leaseTransactions.failed = 1; + secondSource.leaseTransactions.duration = 80; + const idleSource = new QueueSourceStats('https://idle.example.test/queues/jobs'); + const firstQueue = new QueueStats('https://first.example.test/queues/jobs', 'jobs', [ + firstSource, secondSource, idleSource + ]); + + const thirdSource = new QueueSourceStats('https://third.example.test/queues/other'); + thirdSource.leaseTransactions.acquired = 2; + thirdSource.leaseTransactions.contended = 1; + thirdSource.leaseTransactions.failed = 1; + thirdSource.leaseTransactions.duration = 30; + const secondQueue = new QueueStats( + 'https://third.example.test/queues/other', 'other', [thirdSource]); + + const hierarchyStats = new FireleaseStats(() => 0); + hierarchyStats.queues.push(firstQueue, secondQueue); + + assert.strictEqual(firstQueue.leaseTransactions.duration, 50); + assert.strictEqual(secondQueue.leaseTransactions.duration, 30); + assert.strictEqual(hierarchyStats.leaseTransactions.duration, 40); +}); diff --git a/tests/types.ts b/tests/types.ts index 37b0e59..2f11433 100644 --- a/tests/types.ts +++ b/tests/types.ts @@ -1,9 +1,10 @@ import NodeFire from 'nodefire'; import firelease, { RETRY, TESTABLES, attachWorker, blacklist, defaults, extendLease, listTasksInProgress, pingQueues, - settings, shutdown, type Duration, type FireleaseApi, type FireleaseError, - type FireleaseErrorDetails, type FireleaseErrorLevel, type FireleaseSettings, - type FireleaseStats, type Lease, type LeaseItem, type PingReport, type QueueOptions, + settings, shutdown, type CaptureLeaseTransactionMetrics, type Duration, type FireleaseApi, + type FireleaseError, type FireleaseErrorDetails, type FireleaseErrorLevel, type FireleaseSettings, + type FireleaseStats, type Lease, type LeaseItem, type LeaseTransactionStats, type PingReport, + type LeaseTransactionOutcome, type QueueOptions, type QueueMode, type QueueRef, type QueueSourceMode, type QueueSourceStats, type QueueStats, type RetryDirective, type Worker, type WorkerItem, type WorkerResult } from '../src'; @@ -14,6 +15,8 @@ declare const queueRef: NodeFire; declare const generatorWorker: () => Generator; declare const lease: Lease; declare const leaseItem: LeaseItem; +declare const leaseTransactionStats: LeaseTransactionStats; +declare const leaseTransactionOutcome: LeaseTransactionOutcome; declare const workerItem: WorkerItem; declare const fireleaseError: FireleaseError; declare const pingReport: PingReport; @@ -33,17 +36,26 @@ declare const worker: Worker; declare const workerResult: WorkerResult; declare const api: FireleaseApi; void [ - lease, leaseItem, workerItem, fireleaseError, pingReport, queueOptions, duration, errorDetails, - errorLevel, fireleaseSettings, fireleaseStats, queueStats, queueSourceStats, queueMode, + lease, leaseItem, leaseTransactionStats, leaseTransactionOutcome, workerItem, fireleaseError, + pingReport, queueOptions, duration, errorDetails, errorLevel, fireleaseSettings, fireleaseStats, + queueStats, queueSourceStats, queueMode, queueSourceMode, queue, retry, worker, workerResult, api ]; const namedDefaults: QueueOptions = defaults; const namedRetry: RetryDirective = RETRY; const namedSettings: FireleaseSettings = settings; +const captureLeaseTransactionMetrics: CaptureLeaseTransactionMetrics = ( + outcome, tries, transactionDuration +) => { + void [outcome, tries, transactionDuration]; +}; +// @ts-expect-error Lease transaction metric callbacks must be synchronous. +const asyncLeaseTransactionMetrics: CaptureLeaseTransactionMetrics = async () => undefined; void [ TESTABLES, attachWorker, blacklist, extendLease, listTasksInProgress, pingQueues, shutdown, - namedDefaults, namedRetry, namedSettings + namedDefaults, namedRetry, namedSettings, captureLeaseTransactionMetrics, + asyncLeaseTransactionMetrics ]; TESTABLES.resetBetweenTests(); @@ -53,14 +65,21 @@ firelease.attachWorker(queueRef, generatorWorker); firelease.attachWorker(queueRef, item => { const payload = item.payload; const leaseTimeRemaining: number = item.$leaseTimeRemaining; + const firstAcquisition: true | undefined = item._lease.firstAcquisition; void payload; void leaseTimeRemaining; + void firstAcquisition; return firelease.RETRY; }); firelease.attachWorker( [queueRef], - {bufferSize: Infinity, minLease: '30s', preprocess: item => item}, + { + bufferSize: Infinity, + minLease: '30s', + preprocess: item => item, + captureLeaseTransactionMetrics + }, async item => {await firelease.extendLease(item, '1m');} ); @@ -72,12 +91,20 @@ firelease.attachWorker(queueRef, {maxLeaseDelay: '1s'}, () => undefined); firelease.pingQueues(report => { const tasksAcquired: number = report.tasksAcquired; + const acquired: number = report.leaseTransactions.acquired; + const contended: number = report.leaseTransactions.contended; + const failed: number = report.leaseTransactions.failed; + const tries: number = report.leaseTransactions.tries; + const transactionDuration: number = report.leaseTransactions.duration; const sickQueues: (string | null)[] = report.sickQueues; const sickSources: string[] = report.sickSources; const sources: QueueSourceStats[] = report.queues.flatMap(queueResult => queueResult.sources); // @ts-expect-error Lease-delay telemetry was removed with the delay mechanism. void report.leaseDelays; - void [tasksAcquired, sickQueues, sickSources, sources]; + void [ + tasksAcquired, acquired, contended, failed, tries, transactionDuration, sickQueues, sickSources, + sources + ]; }); firelease.settings.globalMaxConcurrent = 10; @@ -91,6 +118,11 @@ const aggregateMode: QueueMode = queueStats.mode; const aggregateSize: number | null = queueStats.size; const aggregateSizeDelta: number | undefined = queueStats.sizeDelta; const aggregateSizeTimestamp: number | undefined = queueStats.sizeTimestamp; +const sourceAcquired: number = queueSourceStats.leaseTransactions.acquired; +const queueContended: number = queueStats.leaseTransactions.contended; +const failedLeaseTransactions: number = queueStats.leaseTransactions.failed; +// @ts-expect-error Lease transaction stats are lifetime totals and cannot be reset. +queueStats.resetLeaseTransactions(); // @ts-expect-error Mutable settings are nested under `settings`. firelease.globalMaxConcurrent = 10; // @ts-expect-error Mutable settings are nested under `settings`. @@ -100,5 +132,5 @@ const taskUrls: string[] = firelease.listTasksInProgress(); void shutdownPromise; void [ taskUrls, currentStats, sizeDelta, aggregateMode, aggregateSize, aggregateSizeDelta, - aggregateSizeTimestamp + aggregateSizeTimestamp, sourceAcquired, queueContended, failedLeaseTransactions ];