import { ILogger } from '../../logger'; import { StreamService } from '../../stream'; import { ThrottleManager } from '../throttling'; import { ErrorHandler } from '../error-handling'; import { LifecycleManager } from '../lifecycle'; import { DuressManager, DuressSnapshot } from '../duress'; import { StreamData, StreamDataResponse } from '../../../types/stream'; import { ProviderClient, ProviderTransaction } from '../../../types/provider'; export declare class ConsumptionManager> { private stream; private logger; private throttleManager; private errorHandler; private lifecycleManager; private reclaimDelay; private reclaimCount; private appId; private role; /** * Consumption stats are written directly to the parent Router so * they are visible in quorum rollcall profiles. */ private get errorCount(); private set errorCount(value); private get counts(); private get hasReachedMaxBackoff(); private set hasReachedMaxBackoff(value); private router; private retry; private duressManager?; private onDuressChange?; private messagesSinceLastEval; private canExtendReservations; private adaptiveReservationTimeout; private adaptiveBatchSize; private lastDepthCheckAt; private static readonly DEPTH_CHECK_INTERVAL_MS; private static readonly DEPTH_SCALE_UP_THRESHOLD; private static readonly DEPTH_SCALE_DOWN_THRESHOLD; private static readonly LEASE_BUFFER_S; constructor(stream: S, logger: ILogger, throttleManager: ThrottleManager, errorHandler: ErrorHandler, lifecycleManager: LifecycleManager, reclaimDelay: number, reclaimCount: number, appId: string, role: any, router: any, retry?: import('../../../types/stream').RetryPolicy, duressManager?: DuressManager); setDuressCallback(callback: (snapshot: DuressSnapshot) => void): void; /** * Adjusts reservation timeout based on stream depth. Called periodically * from the consume loop. When depth is high: * - reservation timeout grows (prevents duplicate re-reservation) * - batch size shrinks (reduces in-memory blocking, shares the stream) * When depth drops, both restore toward configured defaults. */ private adjustConsumptionPressure; createGroup(stream: string, group: string): Promise; publishMessage(topic: string, streamData: StreamData | StreamDataResponse, transaction?: ProviderTransaction): Promise; consumeMessages(stream: string, group: string, consumer: string, callback: (streamData: StreamData) => Promise): Promise; private consumeWithNotifications; private consumeWithPolling; /** * Whether reservation heartbeats are available: the provider must * implement extendReservation and advertise the capability. */ private supportsHeartbeat; consumeOne(stream: string, group: string, id: string, input: StreamData, callback: (streamData: StreamData) => Promise, consumer?: string): Promise; execStreamLeg(input: StreamData, stream: string, id: string, callback: (streamData: StreamData) => Promise): Promise; ackAndDelete(stream: string, group: string, id: string): Promise; ackAndDeleteBatch(stream: string, group: string, ids: string[]): Promise; publishResponse(input: StreamData, output: StreamDataResponse | void): Promise; isStreamMessage(result: any): boolean; }