1#![allow(dead_code)]
29
30pub use backend_optimizer::{
32 BackendOptimizer, BackendPerformance, BackendRecommendation, ConsistencyLevel, CostModel,
33 OptimizationDecision, OptimizationStats, OptimizerConfig, PatternType, WorkloadPattern,
34};
35pub use backpressure::{
36 BackpressureConfig, BackpressureController, BackpressureStats, BackpressureStrategy,
37 FlowControlSignal, RateLimiter as BackpressureRateLimiter,
38};
39pub use bridge::{
40 BridgeInfo, BridgeStatistics, BridgeType, ExternalMessage, ExternalSystemConfig,
41 ExternalSystemType, MessageBridgeManager, MessageTransformer, RoutingRule,
42};
43pub use circuit_breaker::{
44 CircuitBreakerError, CircuitBreakerMetrics, FailureType, SharedCircuitBreakerExt,
45};
46pub use connection_pool::{
47 ConnectionFactory, ConnectionPool, DetailedPoolMetrics, LoadBalancingStrategy, PoolConfig,
48 PoolStats, PoolStatus,
49};
50pub use cqrs::{
51 CQRSConfig, CQRSHealthStatus, CQRSSystem, Command, CommandBus, CommandBusMetrics,
52 CommandHandler, CommandResult, Query, QueryBus, QueryBusMetrics, QueryCacheConfig,
53 QueryHandler, QueryResult as CQRSQueryResult, ReadModelManager, ReadModelMetrics,
54 ReadModelProjection, RetryConfig as CQRSRetryConfig,
55};
56pub use delta::{BatchDeltaProcessor, DeltaComputer, DeltaProcessor, ProcessorStats};
57pub use dlq::{
58 DeadLetterQueue, DlqConfig, DlqEventProcessor, DlqStats as DlqStatsExport, FailedEvent,
59 FailureReason,
60};
61pub use event::{
62 EventCategory, EventMetadata, EventPriority, IsolationLevel, QueryResult as EventQueryResult,
63 SchemaChangeType, SchemaType, SparqlOperationType, StreamEvent,
64};
65pub use event_sourcing::{
66 EventQuery, EventSnapshot, EventStore, EventStoreConfig, PersistenceBackend, QueryOrder,
67 RetentionPolicy, SnapshotConfig, StoredEvent, TimeRange as EventSourcingTimeRange,
68};
69pub use failover::{ConnectionEndpoint, FailoverConfig, FailoverManager};
70pub use graphql_bridge::{
71 BridgeConfig, BridgeStats, GraphQLBridge, GraphQLSubscription, GraphQLUpdate,
72 GraphQLUpdateType, SubscriptionFilter,
73};
74pub use multi_region_replication::{
75 ConflictResolution, ConflictType, GeographicLocation, MultiRegionReplicationManager,
76 RegionConfig, RegionHealth, ReplicatedEvent, ReplicationConfig, ReplicationStats,
77 ReplicationStrategy, VectorClock,
78};
79pub use patch::{PatchParser, PatchSerializer};
80pub use performance_optimizer::{
81 AdaptiveBatcher, AggregationFunction, AutoTuner, BatchPerformancePoint, BatchSizePredictor,
82 BatchingStats, EnhancedMLConfig, MemoryPool, MemoryPoolStats,
83 PerformanceConfig as OptimizerPerformanceConfig, ProcessingResult, ProcessingStats,
84 ProcessingStatus, TuningDecision, ZeroCopyEvent,
85};
86pub use schema_registry::{
87 CompatibilityMode, ExternalRegistryConfig, RegistryAuth, SchemaDefinition, SchemaFormat,
88 SchemaRegistry, SchemaRegistryConfig, ValidationResult, ValidationStats,
89};
90pub use sparql_streaming::{
91 ContinuousQueryManager, QueryManagerConfig, QueryMetadata, QueryResultChannel,
92 QueryResultUpdate, UpdateType,
93};
94pub use store_integration::{
95 ChangeDetectionStrategy, ChangeNotification, RealtimeUpdateManager, StoreChangeDetector,
96 UpdateChannel, UpdateFilter, UpdateNotification,
97};
98
99pub use biological_computing::{
101 AminoAcid, BiologicalProcessingStats, BiologicalStreamProcessor, Cell, CellState,
102 CellularAutomaton, ComputationalFunction, DNASequence, EvolutionaryOptimizer, FunctionalDomain,
103 Individual, Nucleotide, ProteinStructure, SequenceMetadata,
104};
105pub use consciousness_streaming::{
106 ConsciousnessLevel, ConsciousnessStats, ConsciousnessStreamProcessor, DreamSequence,
107 EmotionalContext, IntuitiveEngine, MeditationState,
108};
109pub use disaster_recovery::{
110 BackupCompression, BackupConfig, BackupEncryption, BackupFrequency, BackupJob,
111 BackupRetentionPolicy, BackupSchedule, BackupStatus, BackupStorage, BackupType,
112 BackupVerification, BackupVerificationResult, BackupWindow, BusinessContinuityConfig,
113 ChecksumAlgorithm, CompressionAlgorithm, DRMetrics, DisasterRecoveryConfig,
114 DisasterRecoveryManager, DisasterScenario, EncryptionAlgorithm as BackupEncryptionAlgorithm,
115 FailoverConfig as DRFailoverConfig, ImpactLevel, KeyDerivationFunction, RecoveryConfig,
116 RecoveryOperation, RecoveryPriority, RecoveryRunbook, RecoveryStatus, RecoveryType,
117 ReplicationConfig as DRReplicationConfig, ReplicationMode as DRReplicationMode,
118 ReplicationTarget as DRReplicationTarget, RunbookExecution, RunbookExecutionStatus,
119 RunbookStep, StorageLocation,
120};
121pub use enterprise_audit::{
122 ActionResult, AuditEncryptionConfig, AuditEventType, AuditFilterConfig, AuditMetrics,
123 AuditRetentionConfig, AuditSeverity, AuditStorageBackend, AuditStorageConfig,
124 AuditStreamingConfig, AuthType, ComplianceConfig, ComplianceFinding, ComplianceReport,
125 ComplianceStandard, CompressionType as AuditCompressionType, DestinationAuth, DestinationType,
126 EncryptionAlgorithm, EnterpriseAuditConfig, EnterpriseAuditEvent, EnterpriseAuditLogger,
127 FindingType, KeyManagementConfig, KmsType, S3AuditConfig, StreamingDestination,
128};
129pub use enterprise_monitoring::{
130 Alert, AlertCondition, AlertManager, AlertRule, AlertSeverity as MonitoringAlertSeverity,
131 AlertingConfig, BreachNotificationConfig, ComparisonOperator, EnterpriseMonitoringConfig,
132 EnterpriseMonitoringSystem, EscalationLevel, EscalationPolicy, HealthCheckConfig,
133 HealthCheckEndpoint, HealthCheckType, MeasurementWindow, MetricDefinition, MetricType,
134 MetricValue, MetricsCollector, MetricsConfig, MetricsEndpoint, MetricsEndpointType,
135 MetricsExportConfig, MetricsFormat, NotificationChannel, ProfilingConfig, SlaBreach, SlaConfig,
136 SlaMeasurement, SlaMetricType, SlaObjective, SlaSeverity, SlaStatus, SlaTracker,
137};
138pub use multi_tenancy::{
139 IsolationMode, MultiTenancyConfig, MultiTenancyManager, MultiTenancyMetrics,
140 NamespaceResources, ResourceAllocationStrategy, ResourceType, ResourceUsage, Tenant,
141 TenantLifecycleConfig, TenantNamespace, TenantQuota, TenantStatus, TenantTier,
142};
143pub use observability::{
144 AlertConfig, AlertEvent, AlertSeverity, AlertType, BusinessMetrics, SpanLog, SpanStatus,
145 StreamObservability, StreamingMetrics, TelemetryConfig, TraceSpan,
146};
147pub use performance_utils::{
148 AdaptiveRateLimiter, IntelligentMemoryPool, IntelligentPrefetcher, ParallelStreamProcessor,
149 PerformanceUtilsConfig,
150};
151pub use quantum_communication::{
152 BellState, EntanglementDistribution, QuantumCommConfig, QuantumCommSystem,
153 QuantumOperation as QuantumCommOperation, QuantumSecurityProtocol,
154 QuantumState as QuantumCommState, Qubit,
155};
156pub use quantum_streaming::{
157 QuantumEvent, QuantumOperation, QuantumProcessingStats, QuantumState, QuantumStreamProcessor,
158};
159pub use reliability::{BulkReplayResult, DlqStats, ReplayStatus};
160pub use rsp::{
161 RspConfig, RspLanguage, RspProcessor, RspQuery, StreamClause, StreamDescriptor, Window,
162 WindowConfig, WindowSize, WindowStats, WindowType,
163};
164pub use security::{
165 AuditConfig, AuditLogEntry, AuditLogger, AuthConfig, AuthMethod, AuthenticationProvider,
166 AuthorizationProvider, AuthzConfig, Credentials, EncryptionConfig, Permission, RateLimitConfig,
167 RateLimiter, SecurityConfig as StreamSecurityConfig, SecurityContext, SecurityManager,
168 SecurityMetrics, SessionConfig, ThreatAlert, ThreatDetectionConfig, ThreatDetector,
169};
170pub use temporal_join::{
171 IntervalJoin, JoinResult, LateDataConfig, LateDataStrategy, TemporalJoin, TemporalJoinConfig,
172 TemporalJoinMetrics, TemporalJoinType, TemporalWindow, TimeSemantics, WatermarkConfig,
173 WatermarkStrategy,
174};
175pub use time_travel::{
176 AggregationType, TemporalAggregations, TemporalFilter, TemporalOrdering, TemporalProjection,
177 TemporalQuery, TemporalQueryResult, TemporalResultMetadata, TemporalStatistics, TimePoint,
178 TimeRange as TimeTravelTimeRange, TimeTravelConfig, TimeTravelEngine, TimeTravelMetrics,
179 TimelinePoint,
180};
181pub use tls_security::{
182 CertRotationConfig, CertificateConfig, CertificateFormat, CertificateInfo, CipherSuite,
183 ExpiryWarning, MutualTlsConfig, OcspConfig, RevocationCheckConfig, SessionResumptionConfig,
184 TlsConfig, TlsManager, TlsMetrics, TlsSessionInfo, TlsVersion,
185};
186pub use wasm_edge_computing::{
187 EdgeExecutionResult, EdgeLocation, OptimizationLevel, PerformanceProfile, PluginCapability,
188 PluginSchema, ProcessingSpecialization, ResourceMetrics, SecurityLevel, WasmEdgeConfig,
189 WasmEdgeProcessor, WasmPlugin, WasmProcessingResult, WasmProcessorStats, WasmResourceLimits,
190};
191pub use webhook::{
192 EventFilter as WebhookEventFilter, HttpMethod, RateLimit, RetryConfig as WebhookRetryConfig,
193 WebhookConfig, WebhookInfo, WebhookManager, WebhookMetadata, WebhookSecurity,
194 WebhookStatistics,
195};
196
197pub use custom_serialization::{
199 BenchmarkResults, BsonSerializer, CustomSerializer, FlexBuffersSerializer, IonSerializer,
200 RonSerializer, SerializerBenchmark, SerializerBenchmarkSuite, SerializerRegistry,
201 SerializerStats, ThriftSerializer,
202};
203pub use end_to_end_encryption::{
204 E2EEConfig, E2EEEncryptionAlgorithm, E2EEManager, E2EEStats, EncryptedMessage,
205 HomomorphicEncryption, KeyExchangeAlgorithm, KeyPair, KeyRotationConfig, MultiPartyConfig,
206 ZeroKnowledgeProof,
207};
208#[cfg(feature = "gpu")]
209pub use gpu_acceleration::{
210 AggregationOp, GpuBackend, GpuBuffer, GpuConfig, GpuContext, GpuProcessorConfig, GpuStats,
211 GpuStreamProcessor,
212};
213pub use ml_integration::{
214 AnomalyDetectionAlgorithm, AnomalyDetectionConfig, AnomalyDetector, AnomalyResult,
215 AnomalyStats, FeatureConfig, FeatureExtractor, FeatureVector, MLIntegrationManager,
216 MLModelConfig, ModelMetrics, ModelType, OnlineLearningModel, PredictionResult,
217};
218pub use rate_limiting::{
219 QuotaCheckResult, QuotaEnforcementMode, QuotaLimits, QuotaManager, QuotaOperation,
220 RateLimitAlgorithm, RateLimitConfig as AdvancedRateLimitConfig, RateLimitMonitoringConfig,
221 RateLimitStats as AdvancedRateLimitStats, RateLimiter as AdvancedRateLimiter,
222 RejectionStrategy,
223};
224pub use scalability::{
225 AdaptiveBuffer, AutoScaler, LoadBalancingStrategy as ScalingLoadBalancingStrategy,
226 Node as ScalingNode, NodeHealth, Partition, PartitionManager, PartitionStrategy,
227 ResourceLimits, ResourceUsage as ScalingResourceUsage, ScalingConfig, ScalingDirection,
228 ScalingMode,
229};
230pub use schema_evolution::{
231 CompatibilityCheckResult, CompatibilityIssue, CompatibilityIssueType,
232 CompatibilityMode as SchemaCompatibilityMode, DeprecationInfo, EvolutionResult,
233 FieldDefinition, FieldType, IssueSeverity, MigrationRule, MigrationStrategy, SchemaChange,
234 SchemaDefinition as SchemaEvolutionDefinition, SchemaEvolutionManager,
235 SchemaFormat as SchemaEvolutionFormat, SchemaVersion,
236};
237pub use stream_replay::{
238 EventProcessor, ReplayCheckpoint, ReplayConfig, ReplayFilter, ReplayMode, ReplaySpeed,
239 ReplayStats, ReplayStatus as StreamReplayStatus, ReplayTransformation, StateSnapshot,
240 StreamReplayManager, TransformationType,
241};
242pub use transactional_processing::{
243 IsolationLevel as TransactionalIsolationLevel, LogEntryType, TransactionCheckpoint,
244 TransactionLogEntry, TransactionMetadata, TransactionState, TransactionalConfig,
245 TransactionalProcessor, TransactionalStats,
246};
247pub use zero_copy::{
248 MemoryMappedBuffer, SharedRefBuffer, SimdBatchProcessor, SimdOperation, SplicedBuffer,
249 ZeroCopyBuffer, ZeroCopyConfig, ZeroCopyManager, ZeroCopyStats,
250};
251
252pub use numa_processing::{
254 CpuAffinityMode, HugePageSize, MemoryBandwidthMonitor, MemoryInterleavePolicy, NodeBufferStats,
255 NodeProcessorStats, NumaAllocationStrategy, NumaBuffer, NumaBufferPool, NumaBufferPoolConfig,
256 NumaBufferPoolStats, NumaConfig, NumaNode, NumaProcessorStats, NumaStreamProcessor,
257 NumaThreadPool, NumaThreadPoolStats, NumaTopology, NumaWorker, NumaWorkerStats,
258 WorkerDistributionStrategy,
259};
260pub use out_of_order::{
261 EmitStrategy, GapFillingStrategy, LateEventStrategy, OrderedEvent, OutOfOrderConfig,
262 OutOfOrderHandler, OutOfOrderHandlerBuilder, OutOfOrderStats, SequenceTracker, Watermark,
263};
264pub use performance_profiler::{
265 HistogramStats, LatencyHistogram, OperationTimer, PerformanceProfiler, PerformanceReport,
266 PerformanceSample, PerformanceWarning, ProfilerBuilder, ProfilerConfig, ProfilerStats,
267 Recommendation, RecommendationCategory, RecommendationEffort, RecommendationImpact, Span,
268 WarningSeverity, WarningThresholds, WarningType,
269};
270pub use stream_sql::{
271 AggregateFunction, BinaryOperator, ColumnDefinition, CreateStreamStatement, DataType,
272 Expression, FromClause, JoinType, Lexer, OrderByItem, Parser,
273 QueryResult as StreamSqlQueryResult, QueryType, ResultRow, SelectItem, SelectStatement,
274 SqlValue, StreamMetadata, StreamSqlConfig, StreamSqlEngine, StreamSqlStats, Token,
275 UnaryOperator, WindowSpec, WindowType as SqlWindowType,
276};
277pub use testing_framework::{
278 Assertion, AssertionType, CapturedEvent, EventGenerator, EventMatcher, GeneratorConfig,
279 GeneratorType, MockClock, PerformanceMetric, TestFixture, TestHarness, TestHarnessBuilder,
280 TestHarnessConfig, TestMetrics, TestReport, TestStatus,
281};
282
283pub use anomaly_detection::{
285 Anomaly, AnomalyAlert, AnomalyConfig, AnomalyDetector as AdaptiveAnomalyDetector,
286 AnomalySeverity, AnomalyStats as AdaptiveAnomalyStats, DetectorType, MultiDimensionalDetector,
287};
288pub use migration_tools::{
289 APIMapping, ConceptMapping, GeneratedFile, GeneratedFileType, ManualReviewItem,
290 MigrationConfig, MigrationError, MigrationReport, MigrationSuggestion, MigrationTool,
291 MigrationWarning, QuickStart, ReviewPriority, SourcePlatform, SuggestionCategory,
292};
293pub use online_learning::{
294 ABTestConfig, ABTestResult, DriftDetection, ModelCheckpoint,
295 ModelMetrics as OnlineModelMetrics, ModelType as OnlineModelType, OnlineLearningConfig,
296 OnlineLearningModel as StreamOnlineLearningModel, OnlineLearningStats, Prediction, Sample,
297 StreamFeatureExtractor,
298};
299pub use stream_versioning::{
300 Branch, BranchId, Change, ChangeType, Changeset, Snapshot, StreamVersioning, TimeTravelQuery,
301 TimeTravelTarget, VersionDiff, VersionId, VersionMetadata, VersionedEvent, VersioningConfig,
302 VersioningStats,
303};
304
305pub use automl_stream::{
307 Algorithm, AutoML, AutoMLConfig, AutoMLStats, HyperParameters, ModelPerformance, TaskType,
308 TrainedModel,
309};
310pub use feature_engineering::{
311 Feature, FeatureExtractionConfig, FeatureMetadata, FeaturePipeline, FeatureSet, FeatureStore,
312 FeatureTransform, FeatureValue, ImputationStrategy, PipelineStats,
313};
314pub use neural_architecture_search::{
315 ActivationType, Architecture, ArchitecturePerformance, LayerType, NASConfig, NASStats,
316 ObjectiveWeights, SearchSpace, SearchStrategy, NAS,
317};
318pub use predictive_analytics::{
319 AccuracyMetrics, ForecastAlgorithm, ForecastResult, ForecastingConfig, PredictiveAnalytics,
320 PredictiveStats, SeasonalityType, TrendDirection,
321};
322pub use reinforcement_learning::{
323 Action, Experience, RLAgent, RLAlgorithm, RLConfig, RLStats, State as RLState,
324};
325
326pub use utils::{
328 create_dev_stream, create_prod_stream, BatchProcessor, EventFilter, EventSampler,
329 SimpleRateLimiter, StreamMultiplexer, StreamStats,
330};
331
332pub use advanced_scirs2_optimization::{
334 AdvancedOptimizerConfig, AdvancedStreamOptimizer, MovingStats, OptimizerMetrics,
335};
336pub use cdc_processor::{
337 CdcConfig, CdcConnector, CdcEvent, CdcEventBuilder, CdcMetrics, CdcOperation, CdcProcessor,
338 CdcSource,
339};
340
341pub use adaptive_load_shedding::{
343 DropStrategy, LoadMetrics, LoadSheddingConfig, LoadSheddingManager, LoadSheddingStats,
344};
345
346pub use stream_fusion::{
348 FusableChain, FusedOperation, FusedType, FusionAnalysis, FusionConfig, FusionOptimizer,
349 FusionStats, Operation,
350};
351
352pub use cep_engine::{
354 CepAggregationFunction, CepConfig, CepEngine, CepMetrics, CepStatistics, CompleteMatch,
355 CorrelationFunction, CorrelationResult, CorrelationStats, DetectedPattern, DetectionAlgorithm,
356 DetectionStats, EnrichmentData, EnrichmentService, EnrichmentSource, EnrichmentSourceType,
357 EnrichmentStats, EventBuffer, EventCorrelator, EventPattern, FieldPredicate, PartialMatch,
358 PatternDetector, ProcessingRule, RuleAction, RuleCondition, RuleEngine, RuleExecutionStats,
359 State, StateMachine, TemporalOperator, TimestampedEvent,
360};
361
362pub use data_quality::{
364 AlertCondition as QualityAlertCondition, AlertManager as QualityAlertManager,
365 AlertRule as QualityAlertRule, AlertSeverity as QualityAlertSeverity,
366 AlertStats as QualityAlertStats, AlertType as QualityAlertType, AuditAction, AuditEntry,
367 AuditStats, AuditTrail, CleansingRule, CleansingStats, CorrectionType, DataCleanser,
368 DataCorrection, DataProfiler, DataQualityValidator, DuplicateDetector, DuplicateStats,
369 FailureSeverity, FieldProfile, OutlierMethod, ProfileStats, ProfiledEvent, QualityAlert,
370 QualityConfig, QualityDimension, QualityMetrics, QualityReport, QualityScorer, ScoringStats,
371 ValidationFailure, ValidationResult as QualityValidationResult, ValidationRule,
372};
373
374pub use advanced_sampling::{
376 AdvancedSamplingManager, BloomFilter, BloomFilterStats, CountMinSketch, CountMinSketchStats,
377 HyperLogLog, HyperLogLogStats, ReservoirSampler, ReservoirStats, SamplingConfig,
378 SamplingManagerStats, StratifiedSampler, StratifiedStats, TDigest, TDigestStats,
379};
380
381pub mod backend;
382pub mod backend_optimizer;
383pub mod backpressure;
384pub mod biological_computing;
385pub mod bridge;
386pub mod circuit_breaker;
387pub mod config;
388pub mod connection_pool;
389pub mod connection_pool_health;
390pub mod connection_pool_manager;
391mod connection_pool_tests;
392pub mod connection_pool_types;
393pub mod consciousness_streaming;
394pub mod consumer;
395pub mod cqels;
396pub mod cqrs;
397pub mod csparql;
398pub mod delta;
399pub mod diagnostics;
400pub mod disaster_recovery;
401pub mod dlq;
402pub mod enterprise_audit;
403pub mod enterprise_monitoring;
404pub mod error;
405pub mod event;
406pub mod event_sourcing;
407pub mod failover;
408pub mod graphql_bridge;
409pub mod graphql_subscriptions;
410pub mod health_monitor;
411pub mod join;
412pub mod monitoring;
413pub mod multi_region_replication;
414pub mod multi_tenancy;
415pub mod observability;
416pub mod patch;
417pub mod performance_optimizer;
418pub mod performance_utils;
419pub mod processing;
420pub mod producer;
421pub mod quantum_communication;
422pub mod quantum_processing;
423pub mod quantum_streaming;
424pub mod reconnect;
425pub mod reliability;
426pub mod rsp;
427pub mod schema_registry;
428pub mod security;
429pub mod serialization;
430pub mod serialization_decoder;
431pub mod serialization_encoder;
432pub mod serialization_tests;
433pub mod serialization_types;
434pub mod sparql_streaming;
435pub mod state;
436pub mod store_integration;
437pub mod temporal_join;
438pub mod time_travel;
439pub mod tls_security;
440pub mod types;
441pub mod wasm_edge_computing;
442pub mod webhook;
443
444pub mod custom_serialization;
446pub mod end_to_end_encryption;
447#[cfg(feature = "gpu")]
453pub mod gpu_acceleration;
454pub mod ml_integration;
455pub mod rate_limiting;
456pub mod scalability;
457pub mod schema_evolution;
458pub mod stream_replay;
459pub mod transactional_processing;
460pub mod zero_copy;
461
462pub mod numa_processing;
464pub mod out_of_order;
465pub mod performance_profiler;
466pub mod stream_sql;
467pub mod stream_sql_ast;
468pub mod stream_sql_executor;
469pub mod stream_sql_tests;
470pub mod testing_framework;
471
472pub mod anomaly_detection;
474pub mod migration_tools;
475pub mod online_learning;
476pub mod stream_versioning;
477
478pub mod automl_stream;
480pub mod feature_engineering;
481pub mod neural_architecture_search;
482pub mod predictive_analytics;
483pub mod reinforcement_learning;
484
485pub mod utils;
487
488pub mod advanced_scirs2_optimization;
490pub mod cdc_processor;
491
492pub mod adaptive_load_shedding;
494
495pub mod stream_fusion;
497
498pub mod cep_engine;
500
501pub mod data_quality;
503
504pub mod advanced_sampling;
506
507mod lib_types;
509pub use lib_types::*;
510
511pub mod checkpoint;
513
514pub use checkpoint::{CheckpointCoordinator, CheckpointPhase, GlobalCheckpoint, OperatorSnapshot};
516
517pub use state::{
519 AggregatingState, DeduplicationConfig, DeduplicationLog, DistributedStateBackend,
520 DistributedStateStore, ExactlyOnceProcessor, ExactlyOnceStats, ExactlyOnceTransaction,
521 InMemoryStateBackend, KeyedStateStore, MessageId, PartitionStateValue, StateAggregator,
522 StateBackendStats, StateCoordinator, StatePartition, StatePartitionKey,
523};
524
525pub mod distributed;
527pub mod distributed_state;
528pub mod fault_tolerance;
529
530pub mod websub;
532pub use websub::{
533 DatasetChangeEvent, DatasetEventBus, WebSubHub, WebSubPublisher, WebSubSubscriber,
534};
535
536pub mod ml;
538
539pub use ml::{
541 AnomalyCheckResult, AnomalyDetectorConfig, AnomalyDetectorStats, ExtractedFeatures,
542 FeatureAggregation, FeatureDefinition, FeatureExtractorConfig, ModelConfig, ModelRunnerStats,
543 Prediction as MlPrediction, StreamAnomalyDetector,
544 StreamFeatureExtractor as MlStreamFeatureExtractor, StreamingModelRunner,
545};
546
547pub use distributed_state::manager::{
549 CheckpointConfig as DistributedCheckpointConfig,
550 DeduplicationConfig as DistributedDeduplicationConfig, DeduplicationStats,
551 DistributedStateManager, DistributedStateManagerStats, MigrationPlan, MigrationReason,
552 MigrationStep, OperatorStateSnapshot, PartitionAssignment, SequenceDeduplicator,
553 StateCheckpoint,
554};
555
556pub use fault_tolerance::checkpoint_recovery::{
558 CheckpointManager, CheckpointManagerConfig, CheckpointManagerStats, CheckpointMetadata,
559 PartitionRebalancer, RebalanceAction, RebalancerConfig, RebalancerStats, RecoveryManager,
560 RecoveryManagerStats, RecoveryResult, RecoveryStrategy, StoredCheckpoint, WorkPartition,
561};
562
563pub mod consistency;
565pub mod metrics;
566pub mod watermark;
567
568pub use consistency::ConsistencyManager as StreamConsistencyManager;
570pub use consistency::{
571 ConsistencyConfig, EventualConsistencyBuffer, StreamConsistencyLevel, VersionedValue,
572};
573
574pub use watermark::{
576 LateDataDecision, LateDataHandler, LateDataPolicy, StreamWatermark, WatermarkAligner,
577 WatermarkGenerator,
578};
579
580pub use metrics::{StreamLatencyHistogram, StreamMetrics, StreamMetricsCollector};
582
583pub mod idempotent_delivery;
585
586pub mod window_algebra;
588
589pub mod backpressure_controller;
591
592pub mod sync_schema_registry;
594
595pub mod window_function;
597
598pub use idempotent_delivery::{
601 DeliveryOutcome, HashAlgorithm, IdempotencyKey, IdempotentDeliveryConfig,
602 IdempotentDeliveryManager, IdempotentDeliveryStats, IdempotentProducer, KeyCheckResult,
603};
604
605pub mod stream_checkpoint;
607pub use stream_checkpoint::{Checkpoint, CheckpointStore};
608
609pub mod dead_letter_queue;
611
612pub mod consumer_group;
614
615pub mod stream_router;
617
618pub mod schema_validator;
620
621pub mod message_transformer;
623
624pub mod replay_buffer;
626
627pub mod event_filter;
629
630pub mod visual_designer;
632pub mod visual_designer_engine;
633pub mod visual_designer_tests;
634pub mod visual_designer_types;
635
636pub mod sla;
638pub use sla::{
639 BackpressureAction, SlaBackpressureCoordinator, SlaBackpressureDecision, SlaBackpressurePolicy,
640 StreamAdmissionController, StreamAdmissionDecision, StreamAdmissionStats, StreamSlaConfig,
641};
642
643pub mod window;
647pub use window::{
648 SessionSessionJoin, SessionSessionJoinConfig, TumblingSlidingJoin, TumblingSlidingJoinConfig,
649 TumblingTumblingJoin, TumblingTumblingJoinConfig, WindowJoinKey, WindowJoinResult,
650 WindowJoinStats,
651};
652
653pub mod aggregation;
655pub use aggregation::{
656 ExactlyOnceAggregator, ExactlyOnceAggregatorConfig, ExactlyOnceAggregatorStats,
657 PartitionAggregateState, PartitionAggregateValue,
658};
659
660pub mod neuromorphic_analytics;
662pub mod neuromorphic_analytics_engine;
663pub mod neuromorphic_analytics_learning;
664pub mod neuromorphic_analytics_network;
665pub mod neuromorphic_analytics_patterns;
666mod neuromorphic_analytics_tests;
667pub mod neuromorphic_analytics_types;