import { ILogger } from '../logger'; import { StreamService } from './index'; import { ProviderClient, ProviderTransaction } from '../../types/provider'; import { StreamData, StreamDataResponse } from '../../types/stream'; import { KeyType } from '../../modules/key'; type WorkerCallback = (data: StreamData) => Promise; /** * Process-wide singleton registry that manages one consumer per task queue * (workers) and N consumers per appId (engines). Engine concurrency is * controlled by HMSH_ENGINE_CONCURRENCY — each consumer independently * dequeues from the engine stream using FOR UPDATE SKIP LOCKED. */ declare class StreamConsumerRegistry { private static workerConsumers; private static engineConsumers; /** * Register a worker callback for a (taskQueue, workflowName) pair. * If no consumer exists for this taskQueue, a singleton Router is created. */ static registerWorker(namespace: string, appId: string, guid: string, taskQueue: string, workflowName: string, callback: WorkerCallback, stream: StreamService, store: { mintKey: (type: KeyType, params: any) => string; getThrottleRate: (topic?: string) => Promise; }, logger: ILogger, config?: { reclaimDelay?: number; reclaimCount?: number; readonly?: boolean; retry?: any; }): Promise; /** * Register an engine callback for an appId. * Creates HMSH_ENGINE_CONCURRENCY independent consumers that * dequeue from the engine stream in parallel via SKIP LOCKED. */ static registerEngine(namespace: string, appId: string, guid: string, callback: WorkerCallback, stream: StreamService, store: { mintKey: (type: KeyType, params: any) => string; getThrottleRate: (topic?: string) => Promise; }, logger: ILogger, config?: { reclaimDelay?: number; reclaimCount?: number; }): Promise; /** * Creates a dispatch callback for worker consumers. * Routes messages to the registered callback based on metadata.wfn (workflow_name). */ private static createWorkerDispatcher; /** * Creates a dispatch callback for engine consumers. * Engines are generic processors — the first registered callback handles the message. */ private static createEngineDispatcher; /** * Unregister a worker callback. */ static unregisterWorker(namespace: string, appId: string, taskQueue: string, workflowName: string): Promise; /** * Unregister an engine callback. */ static unregisterEngine(namespace: string, appId: string, callback: WorkerCallback): Promise; /** * Aggregate engine message counts across all consumers for an appId. */ static getEngineCounts(namespace: string, appId: string): { [key: string]: number; }; /** * Stop all consumers and clear the registry. */ static shutdown(): Promise; } export { StreamConsumerRegistry };