import { BehaviorSubject, Observable, Subject } from 'rxjs'; import { Client, debugFnType, IFrame, IMessage, publishParams, StompHeaders } from '@stomp/stompjs'; import { RxStompConfig } from './rx-stomp-config.js'; import { IRxStompPublishParams } from './i-rx-stomp-publish-params.js'; import { RxStompState } from './rx-stomp-state.js'; import { IWatchParams } from './i-watch-params.js'; /** * Main RxJS-friendly STOMP client for browsers and Node.js. * * RxStomp wraps {@link Client} from @stomp/stompjs and exposes key interactions * (connection lifecycle, subscriptions, frames, and errors) as RxJS streams. * * What RxStomp adds: * - Simple, observable-based API for consuming messages. * - Connection lifecycle as BehaviorSubjects/Observables for easy UI binding. * - Transparent reconnection support (configurable through {@link RxStompConfig}). * - Convenience helpers for receipts and server headers. * * Typical lifecycle: * - Instantiate: `const rxStomp = new RxStomp();` * - Configure: `rxStomp.configure({...});` * - Activate: `rxStomp.activate();` * - Consume: `rxStomp.watch({ destination: '/topic/foo' }).subscribe(...)` * - Publish: `rxStomp.publish({ destination: '/topic/foo', body: '...' })` * - Deactivate when done: `await rxStomp.deactivate();` * * Notes: * - Except for `beforeConnect`, all callbacks from @stomp/stompjs are exposed * as RxJS Subjects/Observables here. * - RxStomp tries to transparently handle connection failures. * * Part of `@stomp/rx-stomp`. */ export declare class RxStomp { /** * Connection state as a BehaviorSubject. * * - Emits immediately with the current state (initially `CLOSED`). * - Use this for binding UI widgets or guards based on connection status. * - State values are from {@link RxStompState}. */ readonly connectionState$: BehaviorSubject; /** * Emits whenever a connection is established (including re-connections). * * - If already connected, it emits immediately on subscription. * - The emitted value is always `RxStompState.OPEN`. * - Useful as a trigger to (re)establish subscriptions or flush local queues. */ readonly connected$: Observable; /** * These will be triggered before connectionState$ and connected$. * During reconnecting, it will allow subscriptions to be reinstated before sending * queued messages. */ private _connectionStatePre$; private _connectedPre$; /** * Provides headers from the most recent CONNECTED frame from the broker. * * - Emits on every successful (re)connection. * - If already connected, emits immediately on subscription. * - Typical headers include `server`, `session`, and negotiated `version`. */ readonly serverHeaders$: Observable; protected _serverHeadersBehaviourSubject$: BehaviorSubject; /** * Streams any unhandled MESSAGE frames. * * Use this to receive: * - (RabbitMQ specific) Messages delivered to temporary or auto-named queues. * - Stray messages arriving during unsubscribe processing. * * Emits raw {@link IMessage} instances. * * Maps to: [Client#onUnhandledMessage]{@link Client#onUnhandledMessage}. */ readonly unhandledMessage$: Subject; /** * Streams any unhandled non-MESSAGE, non-ERROR frames. * * Normally unused unless interacting with non-compliant brokers or testing. * * Emits raw {@link IFrame} instances. * * Maps to: [Client#onUnhandledFrame]{@link Client#onUnhandledFrame}. */ readonly unhandledFrame$: Subject