import { ILogger } from '../logger'; import { PublishMessageConfig, StreamConfig, StreamMessage, StreamStats } from '../../types/stream'; import { StringAnyType } from '../../types'; import { KeyStoreParams, KeyType } from '../../modules/key'; import { ProviderClient, ProviderTransaction } from '../../types/provider'; export declare abstract class StreamService { protected streamClient: ClientProvider; protected storeClient: ProviderClient; protected namespace: string; protected logger: ILogger; protected appId: string; protected config: StreamConfig; constructor(streamClient: ClientProvider, storeClient: ProviderClient, config?: StreamConfig); abstract init(namespace: string, appId: string, logger: ILogger): Promise; abstract mintKey(type: KeyType, params: KeyStoreParams): string; abstract createStream(streamName: string): Promise; abstract deleteStream(streamName: string): Promise; abstract createConsumerGroup(streamName: string, groupName: string): Promise; abstract deleteConsumerGroup(streamName: string, groupName: string): Promise; abstract publishMessages(streamName: string, messages: string[], options?: PublishMessageConfig): Promise; abstract consumeMessages(streamName: string, groupName: string, consumerName: string, options?: { batchSize?: number; blockTimeout?: number; autoAck?: boolean; reservationTimeout?: number; enableBackoff?: boolean; initialBackoff?: number; maxBackoff?: number; maxRetries?: number; enableNotifications?: boolean; notificationCallback?: (messages: StreamMessage[]) => void; }): Promise; abstract transact(): ProviderTransaction; abstract ackAndDelete(streamName: string, groupName: string, messageIds: string[], options?: StringAnyType): Promise; abstract acknowledgeMessages(streamName: string, groupName: string, messageIds: string[], options?: StringAnyType): Promise; abstract deleteMessages(streamName: string, groupName: string, messageIds: string[], options?: StringAnyType): Promise; abstract retryMessages(streamName: string, groupName: string, options?: { consumerName?: string; minIdleTime?: number; messageIds?: string[]; delay?: number; maxRetries?: number; limit?: number; }): Promise; reservationTimeout: number; abstract getStreamStats(streamName: string): Promise; /** Whether this instance currently holds the scout role. */ abstract isScout(): boolean; abstract getStreamDepth(streamName: string): Promise; abstract getStreamDepths(streamName: { stream: string; }[]): Promise<{ stream: string; depth: number; }[]>; abstract trimStream(streamName: string, options: { maxLen?: number; maxAge?: number; exactLimit?: boolean; }): Promise; abstract getProviderSpecificFeatures(): { supportsBatching: boolean; supportsDeadLetterQueue: boolean; supportsOrdering: boolean; supportsTrimming: boolean; supportsRetry: boolean; supportsNotifications?: boolean; supportsParallelProcessing?: boolean; supportsReservationExtension?: boolean; maxMessageSize: number; maxBatchSize: number; }; deadLetterMessages?(streamName: string, groupName: string, messageIds: string[]): Promise; expireJobMessages?(jid: string): Promise; extendReservation?(streamName: string, messageId: string, consumerName: string): Promise; stopNotificationConsumer?(streamName: string, groupName: string): Promise; cleanup?(): Promise; }