Skip to main content

PostgresPubSubService

@webda/postgres


Class: PostgresPubSubService<T, K>

Defined in: postgres/src/postgrespubsub.service.ts:72

Pub/sub backed by Postgres' native LISTEN / NOTIFY. A long-lived pg.Client (NOT a pool — pools rotate connections, but each LISTEN is scoped to the connection that issued it) holds the subscription; publishes go through pg_notify(channel, payload). The 8 kB NOTIFY payload cap is enforced in sendMessage — for larger payloads, stash them in a row and notify the row id.

Disconnects trigger a randomized-backoff reconnect so the listener survives transient network or restart blips.

Webda Modda​

PostgresPubSub

Extends​

  • default<T, K>

Type Parameters​

T​

T = any

K​

K extends PostgresPubSubParameters = PostgresPubSubParameters

Constructors​

Constructor​

new PostgresPubSubService<T, K>(name, params): PostgresPubSubService<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​

PostgresPubSubService<T, K>

Inherited from​

PubSubService<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​

PubSubService._compiledCapabilities


[WEBDA_EVENTS]​

[WEBDA_EVENTS]: object

Defined in: core/lib/services/service.d.ts:50

Set the Webda events here

Inherited from​

PubSubService.[WEBDA_EVENTS]


callbacks​

protected callbacks: Set<Subscriber<T>>

Defined in: postgres/src/postgrespubsub.service.ts:92

Local callback registrations. Notifications dispatch to all of them.


client?​

protected optional client?: Client

Defined in: postgres/src/postgrespubsub.service.ts:81

Long-lived listener client. One per service instance because LISTEN is scoped to the connection that issued it — a pool would silently rotate and drop the subscription.


logger​

protected logger: Logger

Defined in: core/lib/services/service.d.ts:59

Logger with class context

Inherited from​

PubSubService.logger


metrics​

protected metrics: 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​

PubSubService.metrics


name​

readonly name: string

Defined in: core/lib/services/iservice.d.ts:19

Inherited from​

PubSubService.name


parameters​

readonly parameters: K

Defined in: core/lib/services/iservice.d.ts:20

Inherited from​

PubSubService.parameters


publishPool?​

protected optional publishPool?: Pool

Defined in: postgres/src/postgrespubsub.service.ts:88

Shared pool used for the publish side (pg_notify). Acquired from the per-config pool registry so multiple services that target the same database share connections instead of each spinning up their own pool. Subscribe-only services skip this and never publish.


stopping​

protected stopping: boolean = false

Defined in: postgres/src/postgrespubsub.service.ts:97

Set during stop so reconnect handlers don't try to come back after teardown.


createConfiguration?​

static optional createConfiguration?: (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​

PubSubService.createConfiguration


filterConfiguration?​

static optional filterConfiguration?: (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​

PubSubService.filterConfiguration


Parameters​

static Parameters: typeof ServiceParameters

Defined in: core/lib/services/service.d.ts:46

Service parameters

Inherited from​

PubSubService.Parameters

Methods​

__clean()​

abstract __clean(): Promise<void>

Defined in: core/lib/services/service.d.ts:215

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

Returns​

Promise<void>

Inherited from​

PubSubService.__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​

PubSubService.addListener


addRoute()​

protected addRoute(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​

PubSubService.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​

PubSubService.authorizeClientEvent


channel()​

protected channel(): string

Defined in: postgres/src/postgrespubsub.service.ts:104

Channel name used for LISTEN/NOTIFY. Resolved at init time so we can default to the service's name when not configured.

Returns​

string

the channel name


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​

PubSubService.computeParameters


connect()​

protected connect(): Promise<void>

Defined in: postgres/src/postgrespubsub.service.ts:121

Open a fresh client, run LISTEN, and wire the notification handler.

Returns​

Promise<void>


consume()​

consume(callback, eventPrototype?, onBind?): CancelablePromise

Defined in: postgres/src/postgrespubsub.service.ts:223

Parameters​

callback​

(event) => Promise<void>

invoked with each event received

eventPrototype?​

() => T

optional class to rehydrate JSON into

onBind?​

() => void

invoked once the subscription is registered

Returns​

CancelablePromise

a cancelable subscription handle

Overrides​

PubSubService.consume


dispatch()​

protected dispatch(payload): void

Defined in: postgres/src/postgrespubsub.service.ts:160

Decode a notification payload and run every registered callback against it.

Parameters​

payload​

string

the raw NOTIFY payload string

Returns​

void


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​

PubSubService.emit


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​

PubSubService.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​

PubSubService.getClientEvents


getMaxListeners()​

getMaxListeners(): number

Defined in: core/lib/events/asynceventemitter.d.ts:84

Returns​

number

Inherited from​

PubSubService.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​

PubSubService.getMetric


getName()​

getName(): string

Defined in: core/lib/services/service.d.ts:209

Get service name

Returns​

string

the result string

Inherited from​

PubSubService.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​

PubSubService.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​

PubSubService.getOperationId


getParameters()​

getParameters(): K

Defined in: core/lib/services/service.d.ts:82

Get the service parameters

Returns​

K

the result

Inherited from​

PubSubService.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​

PubSubService.getService


getState()​

getState(): ServiceStates

Defined in: core/lib/services/service.d.ts:55

Get the current state

Returns​

ServiceStates

the result

Inherited from​

PubSubService.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​

PubSubService.getUrl


handleDisconnect()​

protected handleDisconnect(): void

Defined in: postgres/src/postgrespubsub.service.ts:145

Schedule a reconnect after a short randomized backoff.

Returns​

void


init()​

init(): Promise<PostgresPubSubService<T, K>>

Defined in: postgres/src/postgrespubsub.service.ts:112

Returns​

Promise<PostgresPubSubService<T, K>>

this service

Overrides​

PubSubService.init


initMetrics()​

initMetrics(): void

Defined in: core/lib/queues/pubsubservice.d.ts:20

Returns​

void

Inherited from​

PubSubService.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​

PubSubService.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​

PubSubService.listeners


loadCapabilities()​

protected loadCapabilities(): 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​

PubSubService.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​

PubSubService.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​

PubSubService.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​

PubSubService.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​

PubSubService.once


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​

PubSubService.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​

PubSubService.removeListener


resolve()​

resolve(): this

Defined in: core/lib/services/service.d.ts:98

Resolve parameters and load metadata. @Route-declared routes are scanned and registered by Router.discoverRoutes(services) from Core.init() once every service has resolved — nothing to do per-service here.

Returns​

this

the result

Inherited from​

PubSubService.resolve


sendMessage()​

sendMessage(event): Promise<void>

Defined in: postgres/src/postgrespubsub.service.ts:189

Parameters​

event​

T

the event to publish

Returns​

Promise<void>

Overrides​

PubSubService.sendMessage


setMaxListeners()​

setMaxListeners(n): this

Defined in: core/lib/events/asynceventemitter.d.ts:88

Parameters​

n​

number

Returns​

this

Inherited from​

PubSubService.setMaxListeners


size()​

size(): Promise<number>

Defined in: postgres/src/postgrespubsub.service.ts:212

Returns​

Promise<number>

0 — pub/sub is transient, no queueing

Overrides​

PubSubService.size


stop()​

stop(): Promise<void>

Defined in: postgres/src/postgrespubsub.service.ts:242

Returns​

Promise<void>

Overrides​

PubSubService.stop


toJSON()​

toJSON(): string

Defined in: core/lib/services/service.d.ts:191

Prevent service to be serialized

Returns​

string

the result

Inherited from​

PubSubService.toJSON


toString()​

toString(): string

Defined in: core/lib/services/service.d.ts:91

Return service representation

Returns​

string

the result

Inherited from​

PubSubService.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​

PubSubService.unserialize