Skip to main content

oxirs_stream/
lib.rs

1//! # OxiRS Stream - Ultra-High Performance RDF Streaming Platform
2//!
3//! [![Version](https://img.shields.io/badge/version-0.3.2-blue)](https://github.com/cool-japan/oxirs/releases)
4//! [![docs.rs](https://docs.rs/oxirs-stream/badge.svg)](https://docs.rs/oxirs-stream)
5//!
6//! **Status**: Production Release (v0.3.2)
7//! **Stability**: Public APIs are stable. Production-ready with comprehensive testing.
8//!
9//! Real-time streaming support with Kafka/NATS/Redis I/O, RDF Patch, SPARQL Update delta,
10//! and advanced event processing capabilities.
11//!
12//! This crate provides enterprise-grade real-time data streaming capabilities for RDF datasets,
13//! supporting multiple messaging backends with high-throughput, low-latency guarantees.
14//!
15//! ## Features
16//! - **Multi-Backend Support**: Kafka, NATS JetStream, Redis Streams, AWS Kinesis, Memory
17//! - **High Performance**: 100K+ events/second, <10ms latency, exactly-once delivery
18//! - **Advanced Event Processing**: Real-time pattern detection, windowing, aggregations
19//! - **Enterprise Features**: Circuit breakers, connection pooling, health monitoring
20//! - **Standards Compliance**: RDF Patch protocol, SPARQL Update streaming
21//!
22//! ## Performance Targets
23//! - **Throughput**: 100K+ events/second sustained
24//! - **Latency**: P99 <10ms for real-time processing
25//! - **Reliability**: 99.99% delivery success rate
26//! - **Scalability**: Linear scaling to 1000+ partitions
27
28#![allow(dead_code)]
29
30/// Re-export commonly used types for convenience
31pub 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
99// Stream, StreamConsumer, and StreamProducer are defined below in this module
100pub 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
197// New v0.3.0 feature exports
198pub 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
252// New v0.3.0 exports for developer experience and performance
253pub 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
283// New v0.3.0 exports for ML, versioning, and migration
284pub 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
305// New v0.3.0 advanced ML exports
306pub 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
326// Utility exports
327pub use utils::{
328    create_dev_stream, create_prod_stream, BatchProcessor, EventFilter, EventSampler,
329    SimpleRateLimiter, StreamMultiplexer, StreamStats,
330};
331
332// Advanced SciRS2 optimization exports
333pub 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
341// Adaptive load shedding exports
342pub use adaptive_load_shedding::{
343    DropStrategy, LoadMetrics, LoadSheddingConfig, LoadSheddingManager, LoadSheddingStats,
344};
345
346// Stream fusion optimizer exports
347pub use stream_fusion::{
348    FusableChain, FusedOperation, FusedType, FusionAnalysis, FusionConfig, FusionOptimizer,
349    FusionStats, Operation,
350};
351
352// Complex Event Processing (CEP) engine exports
353pub 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
362// Data quality and validation framework exports
363pub 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
374// Advanced sampling techniques exports
375pub 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
444// New v0.3.0 modules for advanced features
445pub mod custom_serialization;
446pub mod end_to_end_encryption;
447// PRE-EXISTING BUG WORKAROUND (unrelated to the rdkafka/pulsar quarantine):
448// `gpu_acceleration` imports `scirs2_core::gpu`, which scirs2-core gates behind its own
449// `gpu` feature (Pure Rust — `gpu = ["std"]`, no GPU FFI). Gate the module behind an
450// off-by-default `gpu` feature that turns scirs2-core's `gpu` on, so the default build
451// compiles without it and `--all-features` resolves `scirs2_core::gpu`.
452#[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
462// New v0.3.0 modules for developer experience and performance
463pub 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
472// New v0.3.0 modules for ML, versioning, and migration
473pub mod anomaly_detection;
474pub mod migration_tools;
475pub mod online_learning;
476pub mod stream_versioning;
477
478// Advanced ML modules for v0.3.0 completion
479pub mod automl_stream;
480pub mod feature_engineering;
481pub mod neural_architecture_search;
482pub mod predictive_analytics;
483pub mod reinforcement_learning;
484
485// Utilities module
486pub mod utils;
487
488// Advanced SciRS2 optimization module
489pub mod advanced_scirs2_optimization;
490pub mod cdc_processor;
491
492// Adaptive load shedding module
493pub mod adaptive_load_shedding;
494
495// Stream fusion optimizer module
496pub mod stream_fusion;
497
498// Complex Event Processing (CEP) engine module
499pub mod cep_engine;
500
501// Data quality and validation framework module
502pub mod data_quality;
503
504// Advanced sampling techniques module
505pub mod advanced_sampling;
506
507// Extracted type definitions to comply with 2000-line policy
508mod lib_types;
509pub use lib_types::*;
510
511// Distributed stream state management and fault tolerance (v0.2.0)
512pub mod checkpoint;
513
514// Re-export checkpoint types
515pub use checkpoint::{CheckpointCoordinator, CheckpointPhase, GlobalCheckpoint, OperatorSnapshot};
516
517// Re-export distributed state types
518pub 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
525// v0.3.0 modules: Distributed stream processing, distributed state, fault tolerance
526pub mod distributed;
527pub mod distributed_state;
528pub mod fault_tolerance;
529
530// WebSub (W3C) push notifications for RDF dataset changes
531pub mod websub;
532pub use websub::{
533    DatasetChangeEvent, DatasetEventBus, WebSubHub, WebSubPublisher, WebSubSubscriber,
534};
535
536// v1.0.0 ML module for streaming inference and anomaly detection
537pub mod ml;
538
539// Re-export ML types
540pub 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
547// Re-export distributed state manager types (v1.0.0)
548pub 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
556// Re-export fault tolerance checkpoint/recovery types (v1.0.0)
557pub use fault_tolerance::checkpoint_recovery::{
558    CheckpointManager, CheckpointManagerConfig, CheckpointManagerStats, CheckpointMetadata,
559    PartitionRebalancer, RebalanceAction, RebalancerConfig, RebalancerStats, RecoveryManager,
560    RecoveryManagerStats, RecoveryResult, RecoveryStrategy, StoredCheckpoint, WorkPartition,
561};
562
563// v0.2.0 modules: Distributed state, consistency protocols, watermarking, stream metrics
564pub mod consistency;
565pub mod metrics;
566pub mod watermark;
567
568// Re-export consistency types (v0.2.0)
569pub use consistency::ConsistencyManager as StreamConsistencyManager;
570pub use consistency::{
571    ConsistencyConfig, EventualConsistencyBuffer, StreamConsistencyLevel, VersionedValue,
572};
573
574// Re-export watermark types (v0.2.0)
575pub use watermark::{
576    LateDataDecision, LateDataHandler, LateDataPolicy, StreamWatermark, WatermarkAligner,
577    WatermarkGenerator,
578};
579
580// Re-export stream metrics types (v0.2.0)
581pub use metrics::{StreamLatencyHistogram, StreamMetrics, StreamMetricsCollector};
582
583// v1.1.0: Idempotent delivery with idempotency keys
584pub mod idempotent_delivery;
585
586// v1.1.0: Stream Windowing Algebra (tumbling, sliding, session, count-based)
587pub mod window_algebra;
588
589// v1.1.0 round 5: Adaptive backpressure controller (Drop/Block/Throttle/SpillToDisk)
590pub mod backpressure_controller;
591
592// v1.1.0 round 7: Synchronous in-memory schema registry (Kafka Schema Registry style)
593pub mod sync_schema_registry;
594
595// v1.1.0 round 11: Tumbling / sliding / session windowing functions
596pub mod window_function;
597
598// v1.1.0: Event sourcing patterns (append-only log, snapshotting, pub/sub bus)
599// Note: event_sourcing module already declared at line 401
600pub use idempotent_delivery::{
601    DeliveryOutcome, HashAlgorithm, IdempotencyKey, IdempotentDeliveryConfig,
602    IdempotentDeliveryManager, IdempotentDeliveryStats, IdempotentProducer, KeyCheckResult,
603};
604
605// v1.1.0 round 12: Stream checkpoint/offset tracking for at-least-once delivery
606pub mod stream_checkpoint;
607pub use stream_checkpoint::{Checkpoint, CheckpointStore};
608
609// v1.1.0 round 13: Dead letter queue for failed/undeliverable messages
610pub mod dead_letter_queue;
611
612// v1.1.0 round 14: Kafka-style consumer group coordination
613pub mod consumer_group;
614
615// v1.1.0 round 15: Stream message routing (content/topic/header/round-robin/DLQ)
616pub mod stream_router;
617
618// v1.1.0 round 16: Stream message schema validation (field types, formats, strict mode)
619pub mod schema_validator;
620
621// v1.1.0 round 17 (Batch E): Message format transformation pipeline
622pub mod message_transformer;
623
624// v1.1.0 round 18 (Batch E): Event replay buffer with seek and position tracking
625pub mod replay_buffer;
626
627// v1.1.0 round 19: Stream event filtering with composable predicates
628pub mod event_filter;
629
630// Visual pipeline designer and debugger (SVG/JSON/YAML/DOT/Mermaid export)
631pub mod visual_designer;
632pub mod visual_designer_engine;
633pub mod visual_designer_tests;
634pub mod visual_designer_types;
635
636// W2-S6: per-stream SLA admission control + load-shedder coordination.
637pub mod sla;
638pub use sla::{
639    BackpressureAction, SlaBackpressureCoordinator, SlaBackpressureDecision, SlaBackpressurePolicy,
640    StreamAdmissionController, StreamAdmissionDecision, StreamAdmissionStats, StreamSlaConfig,
641};
642
643// W2-S6: watermark-aware window joins (tumbling-tumbling, tumbling-sliding,
644// session-session) — a separate `window` module that complements the
645// time-based `processing::window`.
646pub mod window;
647pub use window::{
648    SessionSessionJoin, SessionSessionJoinConfig, TumblingSlidingJoin, TumblingSlidingJoinConfig,
649    TumblingTumblingJoin, TumblingTumblingJoinConfig, WindowJoinKey, WindowJoinResult,
650    WindowJoinStats,
651};
652
653// W2-S6: exactly-once aggregation under operator parallelism.
654pub mod aggregation;
655pub use aggregation::{
656    ExactlyOnceAggregator, ExactlyOnceAggregatorConfig, ExactlyOnceAggregatorStats,
657    PartitionAggregateState, PartitionAggregateValue,
658};
659
660// Neuromorphic stream analytics (brain-inspired spiking neural networks)
661pub 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;