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 SharedRefBuffer, SimdBatchProcessor, SimdOperation, SplicedBuffer, ZeroCopyBuffer,
249 ZeroCopyConfig, ZeroCopyManager, ZeroCopyStats,
250};
251#[cfg(unix)]
255pub use zero_copy::MemoryMappedBuffer;
256
257pub use numa_processing::{
259 CpuAffinityMode, HugePageSize, MemoryBandwidthMonitor, MemoryInterleavePolicy, NodeBufferStats,
260 NodeProcessorStats, NumaAllocationStrategy, NumaBuffer, NumaBufferPool, NumaBufferPoolConfig,
261 NumaBufferPoolStats, NumaConfig, NumaNode, NumaProcessorStats, NumaStreamProcessor,
262 NumaThreadPool, NumaThreadPoolStats, NumaTopology, NumaWorker, NumaWorkerStats,
263 WorkerDistributionStrategy,
264};
265pub use out_of_order::{
266 EmitStrategy, GapFillingStrategy, LateEventStrategy, OrderedEvent, OutOfOrderConfig,
267 OutOfOrderHandler, OutOfOrderHandlerBuilder, OutOfOrderStats, SequenceTracker, Watermark,
268};
269pub use performance_profiler::{
270 HistogramStats, LatencyHistogram, OperationTimer, PerformanceProfiler, PerformanceReport,
271 PerformanceSample, PerformanceWarning, ProfilerBuilder, ProfilerConfig, ProfilerStats,
272 Recommendation, RecommendationCategory, RecommendationEffort, RecommendationImpact, Span,
273 WarningSeverity, WarningThresholds, WarningType,
274};
275pub use stream_sql::{
276 AggregateFunction, BinaryOperator, ColumnDefinition, CreateStreamStatement, DataType,
277 Expression, FromClause, JoinType, Lexer, OrderByItem, Parser,
278 QueryResult as StreamSqlQueryResult, QueryType, ResultRow, SelectItem, SelectStatement,
279 SqlValue, StreamMetadata, StreamSqlConfig, StreamSqlEngine, StreamSqlStats, Token,
280 UnaryOperator, WindowSpec, WindowType as SqlWindowType,
281};
282pub use testing_framework::{
283 Assertion, AssertionType, CapturedEvent, EventGenerator, EventMatcher, GeneratorConfig,
284 GeneratorType, MockClock, PerformanceMetric, TestFixture, TestHarness, TestHarnessBuilder,
285 TestHarnessConfig, TestMetrics, TestReport, TestStatus,
286};
287
288pub use anomaly_detection::{
290 Anomaly, AnomalyAlert, AnomalyConfig, AnomalyDetector as AdaptiveAnomalyDetector,
291 AnomalySeverity, AnomalyStats as AdaptiveAnomalyStats, DetectorType, MultiDimensionalDetector,
292};
293pub use migration_tools::{
294 APIMapping, ConceptMapping, GeneratedFile, GeneratedFileType, ManualReviewItem,
295 MigrationConfig, MigrationError, MigrationReport, MigrationSuggestion, MigrationTool,
296 MigrationWarning, QuickStart, ReviewPriority, SourcePlatform, SuggestionCategory,
297};
298pub use online_learning::{
299 ABTestConfig, ABTestResult, DriftDetection, ModelCheckpoint,
300 ModelMetrics as OnlineModelMetrics, ModelType as OnlineModelType, OnlineLearningConfig,
301 OnlineLearningModel as StreamOnlineLearningModel, OnlineLearningStats, Prediction, Sample,
302 StreamFeatureExtractor,
303};
304pub use stream_versioning::{
305 Branch, BranchId, Change, ChangeType, Changeset, Snapshot, StreamVersioning, TimeTravelQuery,
306 TimeTravelTarget, VersionDiff, VersionId, VersionMetadata, VersionedEvent, VersioningConfig,
307 VersioningStats,
308};
309
310pub use automl_stream::{
312 Algorithm, AutoML, AutoMLConfig, AutoMLStats, HyperParameters, ModelPerformance, TaskType,
313 TrainedModel,
314};
315pub use feature_engineering::{
316 Feature, FeatureExtractionConfig, FeatureMetadata, FeaturePipeline, FeatureSet, FeatureStore,
317 FeatureTransform, FeatureValue, ImputationStrategy, PipelineStats,
318};
319pub use neural_architecture_search::{
320 ActivationType, Architecture, ArchitecturePerformance, LayerType, NASConfig, NASStats,
321 ObjectiveWeights, SearchSpace, SearchStrategy, NAS,
322};
323pub use predictive_analytics::{
324 AccuracyMetrics, ForecastAlgorithm, ForecastResult, ForecastingConfig, PredictiveAnalytics,
325 PredictiveStats, SeasonalityType, TrendDirection,
326};
327pub use reinforcement_learning::{
328 Action, Experience, RLAgent, RLAlgorithm, RLConfig, RLStats, State as RLState,
329};
330
331pub use utils::{
333 create_dev_stream, create_prod_stream, BatchProcessor, EventFilter, EventSampler,
334 SimpleRateLimiter, StreamMultiplexer, StreamStats,
335};
336
337pub use advanced_scirs2_optimization::{
339 AdvancedOptimizerConfig, AdvancedStreamOptimizer, MovingStats, OptimizerMetrics,
340};
341pub use cdc_processor::{
342 CdcConfig, CdcConnector, CdcEvent, CdcEventBuilder, CdcMetrics, CdcOperation, CdcProcessor,
343 CdcSource,
344};
345
346pub use adaptive_load_shedding::{
348 DropStrategy, LoadMetrics, LoadSheddingConfig, LoadSheddingManager, LoadSheddingStats,
349};
350
351pub use stream_fusion::{
353 FusableChain, FusedOperation, FusedType, FusionAnalysis, FusionConfig, FusionOptimizer,
354 FusionStats, Operation,
355};
356
357pub use cep_engine::{
359 CepAggregationFunction, CepConfig, CepEngine, CepMetrics, CepStatistics, CompleteMatch,
360 CorrelationFunction, CorrelationResult, CorrelationStats, DetectedPattern, DetectionAlgorithm,
361 DetectionStats, EnrichmentData, EnrichmentService, EnrichmentSource, EnrichmentSourceType,
362 EnrichmentStats, EventBuffer, EventCorrelator, EventPattern, FieldPredicate, PartialMatch,
363 PatternDetector, ProcessingRule, RuleAction, RuleCondition, RuleEngine, RuleExecutionStats,
364 State, StateMachine, TemporalOperator, TimestampedEvent,
365};
366
367pub use data_quality::{
369 AlertCondition as QualityAlertCondition, AlertManager as QualityAlertManager,
370 AlertRule as QualityAlertRule, AlertSeverity as QualityAlertSeverity,
371 AlertStats as QualityAlertStats, AlertType as QualityAlertType, AuditAction, AuditEntry,
372 AuditStats, AuditTrail, CleansingRule, CleansingStats, CorrectionType, DataCleanser,
373 DataCorrection, DataProfiler, DataQualityValidator, DuplicateDetector, DuplicateStats,
374 FailureSeverity, FieldProfile, OutlierMethod, ProfileStats, ProfiledEvent, QualityAlert,
375 QualityConfig, QualityDimension, QualityMetrics, QualityReport, QualityScorer, ScoringStats,
376 ValidationFailure, ValidationResult as QualityValidationResult, ValidationRule,
377};
378
379pub use advanced_sampling::{
381 AdvancedSamplingManager, BloomFilter, BloomFilterStats, CountMinSketch, CountMinSketchStats,
382 HyperLogLog, HyperLogLogStats, ReservoirSampler, ReservoirStats, SamplingConfig,
383 SamplingManagerStats, StratifiedSampler, StratifiedStats, TDigest, TDigestStats,
384};
385
386pub mod backend;
387pub mod backend_optimizer;
388pub mod backpressure;
389pub mod biological_computing;
390pub mod bridge;
391pub mod circuit_breaker;
392pub mod config;
393pub mod confluent_registry;
397pub mod connection_pool;
398pub mod connection_pool_health;
399pub mod connection_pool_manager;
400mod connection_pool_tests;
401pub mod connection_pool_types;
402pub mod consciousness_streaming;
403pub mod consumer;
404pub mod cqels;
405pub mod cqrs;
406pub mod csparql;
407pub mod delta;
408pub mod diagnostics;
409pub mod disaster_recovery;
410pub mod dlq;
411pub mod enterprise_audit;
412pub mod enterprise_monitoring;
413pub mod error;
414pub mod event;
415pub mod event_sourcing;
416pub mod failover;
417pub mod graphql_bridge;
418pub mod graphql_subscriptions;
419pub mod health_monitor;
420pub mod join;
421pub mod monitoring;
422pub mod multi_region_replication;
423pub mod multi_tenancy;
424pub mod observability;
425pub mod patch;
426pub mod performance_optimizer;
427pub mod performance_utils;
428pub mod processing;
429pub mod producer;
430pub mod quantum_communication;
431pub mod quantum_processing;
432pub mod quantum_streaming;
433pub mod reconnect;
434pub mod reliability;
435pub mod rsp;
436pub mod schema_registry;
437pub mod security;
438pub mod serialization;
439pub mod serialization_decoder;
440pub mod serialization_encoder;
441pub mod serialization_tests;
442pub mod serialization_types;
443pub mod sparql_streaming;
444pub mod state;
445pub mod store_integration;
446pub mod temporal_join;
447pub mod time_travel;
448pub mod tls_security;
449pub mod types;
450pub mod wasm_edge_computing;
451pub mod wasm_edge_processor;
456pub mod webhook;
457
458pub mod custom_serialization;
460pub mod end_to_end_encryption;
461#[cfg(feature = "gpu")]
467pub mod gpu_acceleration;
468pub mod ml_integration;
469pub mod rate_limiting;
470pub mod scalability;
471pub mod schema_evolution;
472pub mod stream_replay;
473pub mod transactional_processing;
474pub mod zero_copy;
475
476pub mod numa_processing;
478pub mod out_of_order;
479pub mod performance_profiler;
480pub mod stream_sql;
481pub mod stream_sql_ast;
482pub mod stream_sql_executor;
483pub mod stream_sql_tests;
484pub mod testing_framework;
485
486pub mod anomaly_detection;
488pub mod migration_tools;
489pub mod online_learning;
490pub mod stream_versioning;
491
492pub mod automl_stream;
494pub mod feature_engineering;
495pub mod neural_architecture_search;
496pub mod predictive_analytics;
497pub mod reinforcement_learning;
498
499pub mod utils;
501
502pub mod advanced_scirs2_optimization;
504pub mod cdc_processor;
505
506pub mod adaptive_load_shedding;
508
509pub mod stream_fusion;
511
512pub mod cep_engine;
514
515pub mod data_quality;
517
518pub mod advanced_sampling;
520
521mod lib_types;
523pub use lib_types::*;
524
525pub mod checkpoint;
527
528pub use checkpoint::{CheckpointCoordinator, CheckpointPhase, GlobalCheckpoint, OperatorSnapshot};
530
531pub use state::{
533 AggregatingState, DeduplicationConfig, DeduplicationLog, DistributedStateBackend,
534 DistributedStateStore, ExactlyOnceProcessor, ExactlyOnceStats, ExactlyOnceTransaction,
535 InMemoryStateBackend, KeyedStateStore, MessageId, PartitionStateValue, StateAggregator,
536 StateBackendStats, StateCoordinator, StatePartition, StatePartitionKey,
537};
538
539pub mod distributed;
541pub mod distributed_state;
542pub mod fault_tolerance;
543
544pub mod websub;
546pub use websub::{
547 DatasetChangeEvent, DatasetEventBus, WebSubHub, WebSubPublisher, WebSubSubscriber,
548};
549
550pub mod ml;
552
553pub use ml::{
555 AnomalyCheckResult, AnomalyDetectorConfig, AnomalyDetectorStats, ExtractedFeatures,
556 FeatureAggregation, FeatureDefinition, FeatureExtractorConfig, ModelConfig, ModelRunnerStats,
557 Prediction as MlPrediction, StreamAnomalyDetector,
558 StreamFeatureExtractor as MlStreamFeatureExtractor, StreamingModelRunner,
559};
560
561pub use distributed_state::manager::{
563 CheckpointConfig as DistributedCheckpointConfig,
564 DeduplicationConfig as DistributedDeduplicationConfig, DeduplicationStats,
565 DistributedStateManager, DistributedStateManagerStats, MigrationPlan, MigrationReason,
566 MigrationStep, OperatorStateSnapshot, PartitionAssignment, SequenceDeduplicator,
567 StateCheckpoint,
568};
569
570pub use fault_tolerance::checkpoint_recovery::{
572 CheckpointManager, CheckpointManagerConfig, CheckpointManagerStats, CheckpointMetadata,
573 PartitionRebalancer, RebalanceAction, RebalancerConfig, RebalancerStats, RecoveryManager,
574 RecoveryManagerStats, RecoveryResult, RecoveryStrategy, StoredCheckpoint, WorkPartition,
575};
576
577pub mod consistency;
579pub mod metrics;
580pub mod watermark;
581
582pub use consistency::ConsistencyManager as StreamConsistencyManager;
584pub use consistency::{
585 ConsistencyConfig, EventualConsistencyBuffer, StreamConsistencyLevel, VersionedValue,
586};
587
588pub use watermark::{
590 LateDataDecision, LateDataHandler, LateDataPolicy, StreamWatermark, WatermarkAligner,
591 WatermarkGenerator,
592};
593
594pub use metrics::{StreamLatencyHistogram, StreamMetrics, StreamMetricsCollector};
596
597pub mod idempotent_delivery;
599
600pub mod window_algebra;
602
603pub mod backpressure_controller;
605
606pub mod sync_schema_registry;
608
609pub mod window_function;
611
612pub use idempotent_delivery::{
615 DeliveryOutcome, HashAlgorithm, IdempotencyKey, IdempotentDeliveryConfig,
616 IdempotentDeliveryManager, IdempotentDeliveryStats, IdempotentProducer, KeyCheckResult,
617};
618
619pub mod stream_checkpoint;
621pub use stream_checkpoint::{Checkpoint, CheckpointStore};
622
623pub mod dead_letter_queue;
625
626pub mod consumer_group;
628
629pub mod stream_router;
631
632pub mod schema_validator;
634
635pub mod message_transformer;
637
638pub mod replay_buffer;
640
641pub mod event_filter;
643
644pub mod visual_designer;
646pub mod visual_designer_engine;
647pub mod visual_designer_tests;
648pub mod visual_designer_types;
649
650pub mod sla;
652pub use sla::{
653 BackpressureAction, SlaBackpressureCoordinator, SlaBackpressureDecision, SlaBackpressurePolicy,
654 StreamAdmissionController, StreamAdmissionDecision, StreamAdmissionStats, StreamSlaConfig,
655};
656
657pub mod window;
661pub use window::{
662 SessionSessionJoin, SessionSessionJoinConfig, TumblingSlidingJoin, TumblingSlidingJoinConfig,
663 TumblingTumblingJoin, TumblingTumblingJoinConfig, WindowJoinKey, WindowJoinResult,
664 WindowJoinStats,
665};
666
667pub mod aggregation;
669pub use aggregation::{
670 ExactlyOnceAggregator, ExactlyOnceAggregatorConfig, ExactlyOnceAggregatorStats,
671 PartitionAggregateState, PartitionAggregateValue,
672};
673
674pub mod neuromorphic_analytics;
676pub mod neuromorphic_analytics_engine;
677pub mod neuromorphic_analytics_learning;
678pub mod neuromorphic_analytics_network;
679pub mod neuromorphic_analytics_patterns;
680mod neuromorphic_analytics_tests;
681pub mod neuromorphic_analytics_types;