Skip to main content

MemoryQueue

@webda/core


Class: MemoryQueue<T, K>

Defined in: packages/core/src/queues/memoryqueue.service.ts:55

FIFO Queue in Memory

Webda Modda​

Extends​

Type Parameters​

T​

T = any

K​

K extends MemoryQueueParameters = MemoryQueueParameters

Constructors​

Constructor​

new MemoryQueue<T, K>(name, params): MemoryQueue<T, K>

Defined in: packages/core/src/services/service.ts:192

Service

Parameters​

name​

string

The name of the service

params​

K

The parameters block define in the configuration file

Returns​

MemoryQueue<T, K>

Inherited from​

Queue.constructor

Properties​

_compiledCapabilities​

protected _compiledCapabilities: Record<string, any> = {}

Defined in: packages/core/src/services/iservice.ts:58

Capabilities detected at compile-time from @WebdaCapability-tagged interfaces.

Populated during Service.resolve by reading the service's entry in webda.module.json. Each key is a capability name (e.g., "request-filter"), and the value is an empty object {} by default. Override getCapabilities to provide capability-specific configuration or to conditionally disable capabilities.

See​

getCapabilities

Inherited from​

Queue._compiledCapabilities


_timeout​

protected _timeout: Timeout

Defined in: packages/core/src/queues/queueservice.ts:60

Current timeout handler

Inherited from​

Queue._timeout


[WEBDA_EVENTS]​

[WEBDA_EVENTS]: object

Defined in: packages/core/src/services/service.ts:166

Set the Webda events here

Inherited from​

Queue.[WEBDA_EVENTS]


delayer​

protected delayer: WaitDelayer

Defined in: packages/core/src/queues/queueservice.ts:72

Delayer

Inherited from​

Queue.delayer


eventPrototype​

eventPrototype: () => T

Defined in: packages/core/src/queues/queueservice.ts:73

Returns​

T

Inherited from​

Queue.eventPrototype


failedIterations​

protected failedIterations: number

Defined in: packages/core/src/queues/queueservice.ts:68

Current pause instance

Inherited from​

Queue.failedIterations


logger​

protected logger: Logger

Defined in: packages/core/src/services/service.ts:179

Logger with class context

Inherited from​

Queue.logger


metrics​

protected metrics: object

Defined in: packages/core/src/queues/pubsubservice.ts:14

errors​

errors: Counter

messages_pending​

messages_pending: Gauge

messages_received​

messages_received: Counter

messages_sent​

messages_sent: Counter

processing_duration​

processing_duration: Histogram

Inherited from​

Queue.metrics


name​

readonly name: string

Defined in: packages/core/src/services/iservice.ts:35

Inherited from​

Queue.name


parameters​

readonly parameters: K

Defined in: packages/core/src/services/iservice.ts:36

Inherited from​

Queue.parameters


createConfiguration?​

static optional createConfiguration?: (params) => any

Defined in: packages/core/src/services/iservice.ts:41

Create configuration set by the application on load

Parameters​

params​

any

Returns​

any

Inherited from​

Queue.createConfiguration


filterConfiguration?​

static optional filterConfiguration?: (params) => any

Defined in: packages/core/src/services/iservice.ts:46

Create configuration set by the application on load

Parameters​

params​

any

Returns​

any

Inherited from​

Queue.filterConfiguration


Parameters​

static Parameters: typeof ServiceParameters = ServiceParameters

Defined in: packages/core/src/services/service.ts:162

Service parameters

Inherited from​

Queue.Parameters

Methods​

__clean()​

__clean(): Promise<void>

Defined in: packages/core/src/queues/memoryqueue.service.ts:147

Clean the service data, can only be used in test mode

Returns​

Promise<void>

Overrides​

Queue.__clean


addListener()​

addListener<Key>(eventName, listener): this

Defined in: packages/core/src/events/asynceventemitter.ts:83

Type Parameters​

Key​

Key extends never

Parameters​

eventName​

Key

the event name

listener​

(event) => void | Promise<void>

the event listener

Returns​

this

this for chaining

See​

EventEmitter.addListener

Inherited from​

Queue.addListener


addRoute()​

protected addRoute(url, methods, executer, openapi?, override?): void

Defined in: packages/core/src/services/service.ts:384

Add a route dynamicaly

Parameters​

url​

string

of the route can contains dynamic part like {uuid}

methods​

HttpMethodType[]

the HTTP methods

executer​

Function

Method to execute for this route

openapi?​

OpenAPIWebdaDefinition = {}

the OpenAPI specification

override?​

boolean = false

whether to override existing

Returns​

void

Inherited from​

Queue.addRoute


authorizeClientEvent()​

authorizeClientEvent(_event, _context): boolean

Defined in: packages/core/src/services/service.ts:337

Authorize a public event subscription

Parameters​

_event​

string

the event name

_context​

OperationContext

the execution context

Returns​

boolean

true if the condition is met

Inherited from​

Queue.authorizeClientEvent


computeParameters()​

computeParameters(): void

Defined in: packages/core/src/services/service.ts:201

Used to compute or derivate input parameter to attribute

Returns​

void

Deprecated​

Inherited from​

Queue.computeParameters


consume()​

consume(callback, eventPrototype?): CancelablePromise

Defined in: packages/core/src/queues/queueservice.ts:165

Work a queue calling the callback with every Event received If the callback is called without exception the deleteMessage is called

Parameters​

callback​

(event) => Promise<void>

the callback function

eventPrototype?​

() => T

the event prototype

Returns​

CancelablePromise

the result

Inherited from​

Queue.consume


consumerReceiveMessage()​

protected consumerReceiveMessage(): Promise<{ items: number; speed: number; }>

Defined in: packages/core/src/queues/queueservice.ts:101

Receive and process message from the queue

Returns​

Promise<{ items: number; speed: number; }>

the result

Inherited from​

Queue.consumerReceiveMessage


deleteMessage()​

deleteMessage(receipt): Promise<void>

Defined in: packages/core/src/queues/memoryqueue.service.ts:138

Parameters​

receipt​

any

Returns​

Promise<void>

Overrides​

Queue.deleteMessage


emit()​

emit<Key>(event, data): Promise<void>

Defined in: packages/core/src/services/service.ts:473

Emit the event with data and wait for Promise to finish if listener returned a Promise

Type Parameters​

Key​

Key extends never

Parameters​

event​

Key

the event name

data​

object[Key]

the data to process

Returns​

Promise<void>

Inherited from​

Queue.emit


getCapabilities()​

getCapabilities(): Record<string, any>

Defined in: packages/core/src/services/iservice.ts:86

Return the capabilities of this service.

By default returns capabilities detected at compile-time from

Returns​

Record<string, any>

the result

Webda Capability-tagged​

interfaces in webda.module.json.

Override to disable capabilities based on configuration:

getCapabilities() {
const caps = super.getCapabilities();
if (!this.parameters.enabled) delete caps["request-filter"];
return caps;
}

Inherited from​

Queue.getCapabilities


getClientEvents()​

getClientEvents(): string[]

Defined in: packages/core/src/services/service.ts:325

Return the events that an external system can subscribe to

Returns​

string[]

the list of results

Inherited from​

Queue.getClientEvents


getItem()​

getItem<L>(proto?): MessageReceipt<L>[]

Defined in: packages/core/src/queues/memoryqueue.service.ts:93

Type Parameters​

L​

L

Parameters​

proto?​

() => L

Returns​

MessageReceipt<L>[]


getMaxConsumers()​

getMaxConsumers(): number

Defined in: packages/core/src/queues/queueservice.ts:154

Return the max consumers for the queue

It is overridable so if a queue can retrieve several message at once it can just use the worker // and several messages at once

SQS for example will return this.parameters.maxConsumers / 10

Returns​

number

the result number

Inherited from​

Queue.getMaxConsumers


getMaxListeners()​

getMaxListeners(): number

Defined in: packages/core/src/events/asynceventemitter.ts:90

Returns​

number

Inherited from​

Queue.getMaxListeners


getMetric()​

getMetric<T>(type, configuration): T

Defined in: packages/core/src/services/service.ts:310

Add service name label

Type Parameters​

T​

T = Gauge<string> | Counter<string> | Histogram<string>

Parameters​

type​

CustomConstructor<T, [MetricConfiguration<T>]>

the type to look up

configuration​

MetricConfiguration<T>

the configuration

Returns​

T

the result

Inherited from​

Queue.getMetric


getName()​

getName(): string

Defined in: packages/core/src/services/service.ts:483

Get service name

Returns​

string

the result string

Inherited from​

Queue.getName


getOpenApiReplacements()​

getOpenApiReplacements(): any

Defined in: packages/core/src/services/service.ts:413

Return variables for replacement in openapi

Returns​

any

the result

Inherited from​

Queue.getOpenApiReplacements


getOperationId()​

getOperationId(id): string

Defined in: packages/core/src/services/service.ts:371

If undefined is returned it cancel the operation registration

Parameters​

id​

string

the identifier

Returns​

string

the result

Inherited from​

Queue.getOperationId


getParameters()​

getParameters(): K

Defined in: packages/core/src/services/service.ts:209

Get the service parameters

Returns​

K

the result

Inherited from​

Queue.getParameters


getService()​

getService<T>(name): ServicesMap[T]

Defined in: packages/core/src/services/service.ts:300

Get a service by name

Type Parameters​

T​

T extends keyof ServicesMap

Parameters​

name​

T

the name to use

Returns​

ServicesMap[T]

the result map

Deprecated​

Use useService, might reconsider

Inherited from​

Queue.getService


getState()​

getState(): ServiceStates

Defined in: packages/core/src/services/service.ts:172

Get the current state

Returns​

ServiceStates

the result

Inherited from​

Queue.getState


getUrl()​

getUrl(url, _methods): string

Defined in: packages/core/src/services/service.ts:348

Return the full path url based on parameters

Parameters​

url​

string

relative url to service

_methods​

HttpMethodType[]

in case we need filtering (like Store)

Returns​

string

absolute url or undefined if need to skip the Route

Inherited from​

Queue.getUrl


init()​

abstract init(): Promise<MemoryQueue<T, K>>

Defined in: packages/core/src/services/service.ts:463

Will be called after all the Services are created

Returns​

Promise<MemoryQueue<T, K>>

Inherited from​

Queue.init


initMetrics()​

initMetrics(): void

Defined in: packages/core/src/queues/pubsubservice.ts:24

Returns​

void

Inherited from​

Queue.initMetrics


initOperations()​

initOperations(): void

Defined in: packages/core/src/services/service.ts:420

Init the operations from

Returns​

void

Operation​

decorators on this service

Inherited from​

Queue.initOperations


listeners()​

listeners(eventName): Function[]

Defined in: packages/core/src/events/asynceventemitter.ts:183

Get all listeners for an event

Parameters​

eventName​

never

the event name

Returns​

Function[]

the list of results

Inherited from​

Queue.listeners


loadCapabilities()​

protected loadCapabilities(): void

Defined in: packages/core/src/services/service.ts:270

Load capabilities from webda.module.json metadata into _compiledCapabilities.

Called during resolve after dependency injection. Reads the service's type name from parameters, looks it up in the application's module metadata (moddas or beans section), and populates _compiledCapabilities with an empty object for each declared capability name.

Fails silently if the application is not available (e.g., in unit tests where services are instantiated without a full application context).

Returns​

void

Example​

// If webda.module.json contains:
// { "moddas": { "MyApp/HawkService": { "capabilities": ["request-filter", "cors-filter"] } } }
// Then after resolve(), this.getCapabilities() returns:
// { "request-filter": {}, "cors-filter": {} }

See​

getCapabilities

Inherited from​

Queue.loadCapabilities


log()​

log(level, ...args): void

Defined in: packages/core/src/services/service.ts:513

Parameters​

level​

WorkerLogLevel

to log

args​

...any[]

additional arguments

Returns​

void

Inherited from​

Queue.log


off()​

off<Key>(eventName, listener): this

Defined in: packages/core/src/events/asynceventemitter.ts:141

Type Parameters​

Key​

Key extends never

Parameters​

eventName​

Key

the event name

listener​

(event) => void

the event listener

Returns​

this

this for chaining

See​

EventEmitter.off

Inherited from​

Queue.off


on()​

on<Key>(eventName, listener): this

Defined in: packages/core/src/events/asynceventemitter.ts:119

Type Parameters​

Key​

Key extends never

Parameters​

eventName​

Key

the event name

listener​

(event) => void

the event listener

Returns​

this

this for chaining

See​

EventEmitter.once

Inherited from​

Queue.on


once()​

once<Key>(eventName, listener): this

Defined in: packages/core/src/events/asynceventemitter.ts:108

Type Parameters​

Key​

Key extends never

Parameters​

eventName​

Key

the event name

listener​

(event) => void

the event listener

Returns​

this

this for chaining

See​

EventEmitter.once

Inherited from​

Queue.once


receiveMessage()​

receiveMessage<L>(proto?): Promise<MessageReceipt<L>[]>

Defined in: packages/core/src/queues/memoryqueue.service.ts:112

Type Parameters​

L​

L

Parameters​

proto?​

() => L

Returns​

Promise<MessageReceipt<L>[]>

Overrides​

Queue.receiveMessage


removeAllListeners()​

removeAllListeners<Key>(eventName?): this

Defined in: packages/core/src/events/asynceventemitter.ts:150

Type Parameters​

Key​

Key extends never

Parameters​

eventName?​

Key

the event name

Returns​

this

this for chaining

See​

EventEmitter.removeAllListeners

Inherited from​

Queue.removeAllListeners


removeListener()​

removeListener<Key>(eventName, listener): this

Defined in: packages/core/src/events/asynceventemitter.ts:130

Type Parameters​

Key​

Key extends never

Parameters​

eventName​

Key

the event name

listener​

(event) => void

the event listener

Returns​

this

this for chaining

See​

EventEmitter.removeListener

Inherited from​

Queue.removeListener


resolve()​

resolve(): this

Defined in: packages/core/src/queues/queueservice.ts:90

Create the delayer

Returns​

this

this for chaining

Inherited from​

Queue.resolve


sendMessage()​

sendMessage(params): Promise<void>

Defined in: packages/core/src/queues/memoryqueue.service.ts:73

Parameters​

params​

any

Returns​

Promise<void>

Overrides​

Queue.sendMessage


setMaxListeners()​

setMaxListeners(n): this

Defined in: packages/core/src/events/asynceventemitter.ts:97

Parameters​

n​

number

Returns​

this

Inherited from​

Queue.setMaxListeners


size()​

size(): Promise<number>

Defined in: packages/core/src/queues/memoryqueue.service.ts:66

Return queue size

Returns​

Promise<number>

the result number

Overrides​

Queue.size


stop()​

stop(): Promise<void>

Defined in: packages/core/src/services/service.ts:217

Shutdown the current service if action need to be taken

Returns​

Promise<void>

Inherited from​

Queue.stop


toJSON()​

toJSON(): string

Defined in: packages/core/src/services/service.ts:452

Prevent service to be serialized

Returns​

string

the result

Inherited from​

Queue.toJSON


toString()​

toString(): string

Defined in: packages/core/src/services/service.ts:225

Return service representation

Returns​

string

the result

Inherited from​

Queue.toString


unserialize()​

unserialize<L>(data, proto?): L

Defined in: packages/core/src/queues/pubsubservice.ts:59

Unserialize into class

Type Parameters​

L​

L

Parameters​

data​

string

the data to process

proto?​

() => L

the proto

Returns​

L

the result

Inherited from​

Queue.unserialize