import { ILogger } from '../../../logger'; import { PostgresClientType } from '../../../../types/postgres'; import { NotificationConsumer, StreamMessage } from '../../../../types/stream'; import { ProviderClient } from '../../../../types/provider'; /** * Manages PostgreSQL LISTEN/NOTIFY for stream message notifications. * Handles static state shared across all service instances using the same client. * * Channel naming uses table-type prefixes (eng_ / wrk_) instead of group_name, * since engine_streams and worker_streams are separate tables. */ export declare class NotificationManager { private client; private getTableName; private getFallbackInterval; private logger; private static clientNotificationConsumers; private static clientNotificationHandlers; private static clientFallbackPollers; private instanceNotificationConsumers; private notificationHandlerBound; constructor(client: PostgresClientType & ProviderClient, getTableName: () => string, getFallbackInterval: () => number, logger: ILogger); /** * Set up notification handler for this client (once per client). */ setupClientNotificationHandler(serviceInstance: TService): void; /** * Start fallback poller for missed notifications (once per client). */ startClientFallbackPoller(checkForMissedMessages: () => Promise): void; /** * Check for missed messages (fallback polling). */ checkForMissedMessages(fetchMessages: (instance: TService, consumer: NotificationConsumer) => Promise): Promise; /** * Handle incoming PostgreSQL notification. * Channels use table-type prefixes: eng_ for engine, wrk_ for worker. */ private handleNotification; /** * Set up notification consumer for a stream/group. * Uses table-type channel naming (eng_ / wrk_). */ setupNotificationConsumer(serviceInstance: TService, streamName: string, groupName: string, consumerName: string, callback: (messages: StreamMessage[]) => void): Promise; /** * Stop notification consumer for a stream/group. */ stopNotificationConsumer(serviceInstance: TService, streamName: string, groupName: string): Promise; /** * Clean up notification consumers for this instance. */ cleanup(serviceInstance: TService): Promise; /** * Get consumer key from stream and group names. */ private getConsumerKey; } /** * Get configuration values for notification settings. */ export declare function getFallbackInterval(config: any): number; export declare function getNotificationTimeout(config: any): number;