/** * Core interfaces for the NoSQL Time Series Storage Module */ import type { TimeSeriesPoint, TimeSeriesSchema, TimeRange, TimeSeriesQuery, TimeSeriesAggregationQuery, ContinuousQuery, AnomalyDetectionConfig, ForecastConfig, AggregationFunction } from '../../types/timeseries'; export interface TimeSeriesConfig { ingestion: IngestionConfig; storage: StorageConfig; query: QueryConfig; analytics: AnalyticsConfig; realtime: RealtimeConfig; } export interface IngestionConfig { batchSize: number; flushIntervalMs: number; maxConcurrentBatches: number; enableCompression: boolean; enableDeduplication: boolean; writeAheadLog?: { enabled: boolean; maxSizeBytes: number; }; } export interface StorageConfig { hotTierRetentionHours: number; warmTierRetentionDays: number; coldTierRetentionDays: number; compressionEnabled: boolean; compressionType: 'gorilla' | 'zstd' | 'lz4' | 'snappy'; shardingStrategy: 'time-based' | 'metric-based' | 'hash-based'; shardIntervalHours: number; maxShardSizeBytes: number; } export interface QueryConfig { maxPointsPerQuery: number; defaultTimeoutMs: number; cacheEnabled: boolean; cacheTtlMs: number; enableQueryOptimization: boolean; enableParallelExecution: boolean; } export interface AnalyticsConfig { anomalyDetection: { enabled: boolean; defaultSensitivity: number; methods: string[]; maxModelsPerMetric: number; }; forecasting: { enabled: boolean; defaultHorizon: string; methods: string[]; maxHistoryDays: number; }; alerting: { enabled: boolean; maxAlertsPerMinute: number; channels: string[]; }; } export interface RealtimeConfig { maxSubscriptions: number; maxConnectionsPerSubscription: number; heartbeatIntervalMs: number; enableBackpressure: boolean; bufferSizeBytes: number; } export interface TimeSeriesStorageAdapter { name: string; tier: 'hot' | 'warm' | 'cold'; write(points: TimeSeriesPoint[]): Promise; read(query: TimeSeriesQuery): Promise; delete(metric: string, timeRange: TimeRange): Promise; writeBatch(batches: TimeSeriesPoint[][]): Promise; readBatch(queries: TimeSeriesQuery[]): Promise; createSchema(schema: TimeSeriesSchema): Promise; updateSchema(metric: string, schema: Partial): Promise; getSchema(metric: string): Promise; listSchemas(): Promise; compact(metric?: string): Promise; vacuum(): Promise; getStats(): Promise; } export interface StorageStats { totalPoints: number; totalMetrics: number; totalSizeBytes: number; compressionRatio: number; queryLatencyMs: { p50: number; p95: number; p99: number; }; writeThroughput: { pointsPerSecond: number; bytesPerSecond: number; }; } export interface IngestionEngine { writePoint(point: TimeSeriesPoint): Promise; writeBatch(points: TimeSeriesPoint[]): Promise; createWriteStream(): WritableStream; writePrometheusMetrics(metrics: string): Promise; writeInfluxDBLine(line: string): Promise; writeStatsD(metric: string): Promise; createIngestionStream(options?: StreamOptions): Promise>; getIngestionStats(): Promise; } export interface StreamOptions { batchSize?: number; flushInterval?: number; compression?: boolean; deduplication?: boolean; onError?: (error: Error) => void; onSuccess?: (batchSize: number) => void; } export interface IngestionStats { pointsIngested: number; batchesProcessed: number; errorsCount: number; duplicatesDropped: number; avgBatchSize: number; ingestionRatePerSecond: number; bufferUtilization: number; } export interface QueryEngine { query(query: TimeSeriesQuery): Promise; aggregate(query: TimeSeriesAggregationQuery): Promise; movingAverage(options: MovingAverageOptions): Promise; derivative(metric: string, timeRange: TimeRange): Promise; histogram(metric: string, timeRange: TimeRange, buckets: number): Promise; timeBucket(interval: string, timestamp: number): number; timeGap(points: TimeSeriesPoint[], maxGapMs: number): TimeGap[]; interpolate(points: TimeSeriesPoint[], method: 'linear' | 'spline'): TimeSeriesPoint[]; createLiveQuery(query: TimeSeriesQuery): Promise>; createContinuousQuery(cq: ContinuousQuery): Promise; explainQuery(query: TimeSeriesQuery): Promise; optimizeQuery(query: TimeSeriesQuery): Promise; } export interface TimeSeriesQueryResult { points: TimeSeriesPoint[]; metadata: { totalPoints: number; executionTimeMs: number; fromCache: boolean; scannedShards: string[]; }; } export interface TimeSeriesAggregationResult { buckets: Array<{ timestamp: number; values: Record; groupBy?: Record; }>; metadata: { totalBuckets: number; executionTimeMs: number; fromCache: boolean; }; } export interface MovingAverageOptions { metric: string; timeRange: TimeRange; window: string; step?: string; } export interface TimeGap { start: number; end: number; durationMs: number; } export interface HistogramResult { buckets: Array<{ min: number; max: number; count: number; }>; stats: { total: number; min: number; max: number; avg: number; }; } export interface QueryPlan { steps: QueryStep[]; estimatedCost: number; estimatedRows: number; tiersAccessed: string[]; indexesUsed: string[]; } export interface QueryStep { operation: string; description: string; cost: number; parallelizable: boolean; } export interface AnomalyDetector { trainModel(metric: string, config: AnomalyDetectionConfig): Promise; updateModel(modelId: string, data: TimeSeriesPoint[]): Promise; deleteModel(modelId: string): Promise; detectAnomalies(metric: string, timeRange: TimeRange, config?: AnomalyDetectionConfig): Promise; detectRealtime(point: TimeSeriesPoint, config: AnomalyDetectionConfig): Promise; detectBatch(points: TimeSeriesPoint[], config: AnomalyDetectionConfig): Promise; } export interface Anomaly { timestamp: number; value: number; expected: number; score: number; tags: Record; method: string; severity: 'low' | 'medium' | 'high' | 'critical'; } export interface ForecastingEngine { forecast(config: ForecastConfig): Promise; forecastRealtime(metric: string, horizon: string): Promise>; trainForecastModel(metric: string, config: ForecastConfig): Promise; evaluateModel(modelId: string): Promise; } export interface ForecastResult { timestamp: number; value: number; upperBound?: number; lowerBound?: number; confidence?: number; } export interface ModelEvaluation { mae: number; mse: number; rmse: number; mape: number; r2: number; lastUpdated: string; } export interface AlertingSystem { createAlertRule(rule: AlertRule): Promise; updateAlertRule(ruleId: string, rule: Partial): Promise; deleteAlertRule(ruleId: string): Promise; getAlertRule(ruleId: string): Promise; listAlertRules(): Promise; evaluateAlerts(): Promise; evaluateAlert(ruleId: string): Promise; fireAlert(alert: Alert): Promise; acknowledgeAlert(alertId: string, user: string): Promise; resolveAlert(alertId: string, user: string): Promise; addNotificationChannel(channel: NotificationChannel): Promise; removeNotificationChannel(channelId: string): Promise; } export interface AlertRule { id: string; name: string; metric: string; condition: AlertCondition; threshold: number; evaluationInterval: string; forDuration?: string; labels: Record; annotations: Record; notificationChannels: string[]; silenceRules?: SilenceRule[]; enabled: boolean; } export interface AlertCondition { operator: '>' | '<' | '>=' | '<=' | '==' | '!='; aggregation: AggregationFunction; timeRange: string; groupBy?: string[]; } export interface SilenceRule { matchers: Record; startTime: string; endTime: string; createdBy: string; comment: string; } export interface Alert { id: string; ruleId: string; metric: string; value: number; threshold: number; condition: string; severity: 'info' | 'warning' | 'critical'; status: 'firing' | 'resolved' | 'acknowledged'; startsAt: string; endsAt?: string; acknowledgedBy?: string; acknowledgedAt?: string; labels: Record; annotations: Record; } export interface NotificationChannel { id: string; type: 'email' | 'webhook' | 'slack' | 'pagerduty' | 'sms'; name: string; config: Record; enabled: boolean; } export interface CompressionManager { compress(data: TimeSeriesPoint[], algorithm: string): Promise; decompress(block: CompressedBlock): Promise; gorillaCompress(values: number[], timestamps: number[]): Promise; gorillaDecompress(data: Uint8Array): Promise<{ values: number[]; timestamps: number[]; }>; zstdCompress(data: Uint8Array, level?: number): Promise; zstdDecompress(data: Uint8Array): Promise; getCompressionStats(): Promise; } export interface CompressedBlock { algorithm: string; originalSize: number; compressedSize: number; compressionRatio: number; data: Uint8Array; checksum: string; metadata: { pointCount: number; startTime: number; endTime: number; metrics: string[]; }; } export interface CompressionStats { totalBlocks: number; totalOriginalBytes: number; totalCompressedBytes: number; averageCompressionRatio: number; algorithmStats: Record; } export interface ProtocolAdapter { name: string; version: string; parse(data: string | Uint8Array): Promise; format(points: TimeSeriesPoint[]): Promise; validate(data: string | Uint8Array): Promise; getSchema?(): Promise; } export interface RealtimeManager { subscribe(subscription: TimeSeriesSubscription): Promise; unsubscribe(subscriptionId: string): Promise; broadcast(points: TimeSeriesPoint[]): Promise; broadcastToSubscription(subscriptionId: string, points: TimeSeriesPoint[]): Promise; addConnection(connectionId: string, metadata?: Record): Promise; removeConnection(connectionId: string): Promise; getActiveConnections(): Promise; pauseSubscription(subscriptionId: string): Promise; resumeSubscription(subscriptionId: string): Promise; } export interface TimeSeriesSubscription { id: string; connectionId: string; metrics: string[]; filters?: Record; aggregateWindow?: string; maxPointsPerSecond?: number; bufferSize?: number; } //# sourceMappingURL=interfaces.d.ts.map