Skip to main content

SocketPubSubService

@webda/fs


Class: SocketPubSubService<T, K>

Defined in: fs/src/socketpubsub.service.ts:62

Pub/sub backed by a unix-domain socket. Designed for single-host IPC: one peer binds the socket file and acts as broker, fanning each published message out to every connected client; other peers connect as clients. Any peer can publish — non-broker publishers send through the socket and the broker fanouts. Every connected peer (publisher included) receives every message, so the in-process semantics match the cross-process ones.

If the broker disconnects, surviving clients reconnect after a short randomized delay; one of them wins the bind race and becomes the new broker. Stale socket files left by killed brokers are detected via connect() returning ECONNREFUSED and unlinked before re-binding.

Wire format: 4-byte big-endian uint32 length prefix followed by a UTF-8 JSON payload.

Webda Modda​

SocketPubSub

Extends​

  • default<T, K>

Type Parameters​

T​

T = any

K​

K extends SocketPubSubParameters = SocketPubSubParameters

Constructors​

Constructor​

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

SocketPubSubService<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: fs/src/socketpubsub.service.ts:84

Local callback registrations made via consume. The broker dispatches to these directly; client-mode peers feed them from the broker socket's incoming frames.


clientSocket?​

protected optional clientSocket?: Socket

Defined in: fs/src/socketpubsub.service.ts:74

Outbound socket to the broker, when this peer is a client.


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


server?​

protected optional server?: Server

Defined in: fs/src/socketpubsub.service.ts:70

Broker-mode listener. Set when this peer is the broker, undefined otherwise. Mutually exclusive with clientSocket.


stopping​

protected stopping: boolean = false

Defined in: fs/src/socketpubsub.service.ts:89

Set during stop so disconnect handlers don't try to reconnect after we've torn down.


subscribers​

protected subscribers: Set<Socket>

Defined in: fs/src/socketpubsub.service.ts:78

Active subscriber sockets the broker fans messages out to.


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


bindAsBroker()​

protected bindAsBroker(path): Promise<void>

Defined in: fs/src/socketpubsub.service.ts:141

Bind the socket file and start listening. On every accepted connection we frame-decode incoming bytes, fanout the raw frame to all OTHER subscribers, and dispatch the parsed payload to local callbacks.

Parameters​

path​

string

filesystem path of the unix socket

Returns​

Promise<void>

resolves once the listener is up


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: fs/src/socketpubsub.service.ts:105

Connect to an existing broker, or bind ourselves as broker if there isn't one. Handles the EADDRINUSE race by retrying as a client.

Returns​

Promise<void>


connectAsClient()​

protected connectAsClient(path): Promise<void>

Defined in: fs/src/socketpubsub.service.ts:169

Open a client socket to the broker.

Parameters​

path​

string

filesystem path of the unix socket

Returns​

Promise<void>

resolves once connected, rejects on the first connection error


consume()​

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

Defined in: fs/src/socketpubsub.service.ts:288

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


dispatchToCallbacks()​

protected dispatchToCallbacks(buf): void

Defined in: fs/src/socketpubsub.service.ts:227

Decode a JSON frame and run every registered callback against it. Errors thrown by callbacks are swallowed (logged + counted) so one bad subscriber doesn't break others.

Parameters​

buf​

Buffer

the unframed JSON payload bytes

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


fanout()​

protected fanout(buf, except?): void

Defined in: fs/src/socketpubsub.service.ts:213

Write a framed payload to every subscriber except the optional sender (used to avoid echoing a client's own message back to it via the broker — the client receives it via its local callback dispatch instead).

Parameters​

buf​

Buffer

the unframed JSON payload bytes

except?​

Socket

optional socket to exclude from fanout

Returns​

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​

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

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


handleBrokerDisconnect()​

protected handleBrokerDisconnect(): void

Defined in: fs/src/socketpubsub.service.ts:195

Schedule a reconnect attempt after a randomized backoff. The jitter keeps surviving clients from all racing to bind the new socket on the same tick.

Returns​

void


init()​

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

Defined in: fs/src/socketpubsub.service.ts:95

Returns​

Promise<SocketPubSubService<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: fs/src/socketpubsub.service.ts:256

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: fs/src/socketpubsub.service.ts:277

Returns​

Promise<number>

0 — pub/sub is transient, no queueing

Overrides​

PubSubService.size


stop()​

stop(): Promise<void>

Defined in: fs/src/socketpubsub.service.ts:312

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