PostgresQueueService
Class: PostgresQueueService<T, K>
Defined in: postgres/src/postgresqueue.service.ts:88
Postgres-backed FIFO queue using SELECT … FOR UPDATE SKIP LOCKED (PG
9.5+) for atomic multi-worker pulls. A schema-managed table holds
pending and locked rows; receive locks a batch atomically, delete (or
the visibility-timeout sweep) clears them. No extra infrastructure
needed beyond a Postgres connection — reuses the same DB you're
already running for the store.
Wire format: payload column is jsonb so messages survive
round-tripping with their structure intact and can be queried directly
if you ever need to inspect the queue.
Webda Modda
PostgresQueue
Extends
Queue<T,K>
Type Parameters
T
T = any
K
K extends PostgresQueueParameters = PostgresQueueParameters
Constructors
Constructor
new PostgresQueueService<
T,K>(name,params):PostgresQueueService<T,K>
Defined in: core/lib/services/service.d.ts:72
Service
Parameters
name
string
The name of the service
params
K
The parameters block define in the configuration file
Returns
PostgresQueueService<T, K>
Inherited from
Queue<T, K>.constructor
Properties
_compiledCapabilities
protected_compiledCapabilities:Record<string,any>
Defined in: core/lib/services/iservice.d.ts:39
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: core/lib/queues/queueservice.d.ts:51
Current timeout handler
Inherited from
Queue._timeout
[WEBDA_EVENTS]
[WEBDA_EVENTS]:
object
Defined in: core/lib/services/service.d.ts:50
Set the Webda events here
Inherited from
Queue.[WEBDA_EVENTS]
client?
protectedoptionalclient?:Client|Pool
Defined in: postgres/src/postgresqueue.service.ts:96
Backing pg client or pool. Pools are preferred under load — receive locks rotate across connections and benefit from concurrency.
delayer
protecteddelayer:WaitDelayer
Defined in: core/lib/queues/queueservice.d.ts:63
Delayer
Inherited from
Queue.delayer
eventPrototype
eventPrototype: () =>
T
Defined in: core/lib/queues/queueservice.d.ts:64
Returns
T
Inherited from
Queue.eventPrototype
failedIterations
protectedfailedIterations:number
Defined in: core/lib/queues/queueservice.d.ts:59
Current pause instance
Inherited from
Queue.failedIterations
logger
protectedlogger:Logger
Defined in: core/lib/services/service.d.ts:59
Logger with class context
Inherited from
Queue.logger
metrics
protectedmetrics:object
Defined in: core/lib/queues/pubsubservice.d.ts:10
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
readonlyname:string
Defined in: core/lib/services/iservice.d.ts:19
Inherited from
Queue.name
ownsPool
protectedownsPool:boolean=false
Defined in: postgres/src/postgresqueue.service.ts:101
True when client is a shared pool acquired from the registry. Drives the matching releasePool call in stop.
parameters
readonlyparameters:K
Defined in: core/lib/services/iservice.d.ts:20
Inherited from
Queue.parameters
createConfiguration?
staticoptionalcreateConfiguration?: (params) =>any
Defined in: core/lib/services/iservice.d.ts:24
Create configuration set by the application on load
Parameters
params
any
Returns
any
Inherited from
Queue.createConfiguration
filterConfiguration?
staticoptionalfilterConfiguration?: (params) =>any
Defined in: core/lib/services/iservice.d.ts:28
Create configuration set by the application on load
Parameters
params
any
Returns
any
Inherited from
Queue.filterConfiguration
Parameters
staticParameters: typeofServiceParameters
Defined in: core/lib/services/service.d.ts:46
Service parameters
Inherited from
Queue.Parameters
Accessors
table
Get Signature
get
protectedtable():string
Defined in: postgres/src/postgresqueue.service.ts:108
Resolved table identifier. Validated at init to keep the table name out of any literal SQL paths that aren't parameterizable.
Returns
string
the table name
Methods
__clean()
__clean():
Promise<void>
Defined in: postgres/src/postgresqueue.service.ts:241
Convenience: drop and recreate the queue table. Used by tests; not
for production. Mirrors the __clean hook on FileQueue.
Returns
Promise<void>
Overrides
Queue.__clean
addListener()
addListener<
Key>(eventName,listener):this
Defined in: core/lib/events/asynceventemitter.d.ts:80
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()
protectedaddRoute(url,methods,executer,openapi?,override?):void
Defined in: core/lib/services/service.d.ts:177
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
whether to override existing
Returns
void
Inherited from
Queue.addRoute
authorizeClientEvent()
authorizeClientEvent(
_event,_context):boolean
Defined in: core/lib/services/service.d.ts:153
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: core/lib/services/service.d.ts:77
Used to compute or derivate input parameter to attribute
Returns
void
Deprecated
Inherited from
Queue.computeParameters
consume()
consume(
callback,eventPrototype?):CancelablePromise
Defined in: core/lib/queues/queueservice.d.ts:108
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()
protectedconsumerReceiveMessage():Promise<{items:number;speed:number; }>
Defined in: core/lib/queues/queueservice.d.ts:86
Receive and process message from the queue
Returns
Promise<{ items: number; speed: number; }>
the result
Inherited from
Queue.consumerReceiveMessage
deleteMessage()
deleteMessage(
id):Promise<void>
Defined in: postgres/src/postgresqueue.service.ts:220
Parameters
id
string
the receipt handle returned by receiveMessage
Returns
Promise<void>
Overrides
Queue.deleteMessage
emit()
emit<
Key>(event,data):Promise<void>
Defined in: core/lib/services/service.d.ts:204
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
ensureTable()
protectedensureTable():Promise<void>
Defined in: postgres/src/postgresqueue.service.ts:138
Create the queue table and the index that supports the SKIP LOCKED receive query, if they're missing.
Returns
Promise<void>
getCapabilities()
getCapabilities():
Record<string,any>
Defined in: core/lib/services/iservice.d.ts:61
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: core/lib/services/service.d.ts:144
Return the events that an external system can subscribe to
Returns
string[]
the list of results
Inherited from
Queue.getClientEvents
getMaxConsumers()
getMaxConsumers():
number
Defined in: postgres/src/postgresqueue.service.ts:233
Override the queue's per-receive parallelism: receiveMessage already
pulls a batch, so the consumer-spawning loop only needs one parent
worker per batchSize.
Returns
number
the result number
Overrides
Queue.getMaxConsumers
getMaxListeners()
getMaxListeners():
number
Defined in: core/lib/events/asynceventemitter.d.ts:84
Returns
number
Inherited from
Queue.getMaxListeners
getMetric()
getMetric<
T>(type,configuration):T
Defined in: core/lib/services/service.d.ts:138
Add service name label
Type Parameters
T
T = Counter<string> | Histogram<string> | Gauge<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: core/lib/services/service.d.ts:209
Get service name
Returns
string
the result string
Inherited from
Queue.getName
getOpenApiReplacements()
getOpenApiReplacements():
any
Defined in: core/lib/services/service.d.ts:182
Return variables for replacement in openapi
Returns
any
the result
Inherited from
Queue.getOpenApiReplacements
getOperationId()
getOperationId(
id):string
Defined in: core/lib/services/service.d.ts:167
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: core/lib/services/service.d.ts:82
Get the service parameters
Returns
K
the result
Inherited from
Queue.getParameters
getService()
getService<
T>(name):ServicesMap[T]
Defined in: core/lib/services/service.d.ts:131
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: core/lib/services/service.d.ts:55
Get the current state
Returns
ServiceStates
the result
Inherited from
Queue.getState
getUrl()
getUrl(
url,_methods):string
Defined in: core/lib/services/service.d.ts:161
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()
init():
Promise<PostgresQueueService<T,K>>
Defined in: postgres/src/postgresqueue.service.ts:116
Returns
Promise<PostgresQueueService<T, K>>
this service
Overrides
Queue.init
initMetrics()
initMetrics():
void
Defined in: core/lib/queues/pubsubservice.d.ts:20
Returns
void
Inherited from
Queue.initMetrics
initOperations()
initOperations():
void
Defined in: core/lib/services/service.d.ts:186
Init the operations from
Returns
void
Operation
decorators on this service
Inherited from
Queue.initOperations
listeners()
listeners(
eventName):Function[]
Defined in: core/lib/events/asynceventemitter.d.ts:139
Get all listeners for an event
Parameters
eventName
never
the event name
Returns
Function[]
the list of results
Inherited from
Queue.listeners
loadCapabilities()
protectedloadCapabilities():void
Defined in: core/lib/services/service.d.ts:120
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: core/lib/services/service.d.ts:226
Parameters
level
WorkerLogLevel
to log
args
...any[]
additional arguments
Returns
void
Inherited from
Queue.log
off()
off<
Key>(eventName,listener):this
Defined in: core/lib/events/asynceventemitter.d.ts:116
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: core/lib/events/asynceventemitter.d.ts:102
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: core/lib/events/asynceventemitter.d.ts:95
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: postgres/src/postgresqueue.service.ts:187
Type Parameters
L
L
Parameters
proto?
() => L
optional prototype to rehydrate the payload into
Returns
Promise<MessageReceipt<L>[]>
the locked batch
Overrides
Queue.receiveMessage
removeAllListeners()
removeAllListeners<
Key>(eventName?):this
Defined in: core/lib/events/asynceventemitter.d.ts:122
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: core/lib/events/asynceventemitter.d.ts:109
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: core/lib/queues/queueservice.d.ts:80
Create the delayer
Returns
this
this for chaining
Inherited from
Queue.resolve
sendMessage()
sendMessage(
event):Promise<void>
Defined in: postgres/src/postgresqueue.service.ts:164
Parameters
event
T
the event to enqueue
Returns
Promise<void>
Overrides
Queue.sendMessage
setMaxListeners()
setMaxListeners(
n):this
Defined in: core/lib/events/asynceventemitter.d.ts:88
Parameters
n
number
Returns
this
Inherited from
Queue.setMaxListeners
size()
size():
Promise<number>
Defined in: postgres/src/postgresqueue.service.ts:174
Returns
Promise<number>
count of messages currently visible (pending or expired-lock)
Overrides
Queue.size
stop()
stop():
Promise<void>
Defined in: postgres/src/postgresqueue.service.ts:253
Returns
Promise<void>
Overrides
Queue.stop
toJSON()
toJSON():
string
Defined in: core/lib/services/service.d.ts:191
Prevent service to be serialized
Returns
string
the result
Inherited from
Queue.toJSON
toString()
toString():
string
Defined in: core/lib/services/service.d.ts:91
Return service representation
Returns
string
the result
Inherited from
Queue.toString
unserialize()
unserialize<
L>(data,proto?):L
Defined in: core/lib/queues/pubsubservice.d.ts:32
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