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.4.1-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.4.1)
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    SharedRefBuffer, SimdBatchProcessor, SimdOperation, SplicedBuffer, ZeroCopyBuffer,
249    ZeroCopyConfig, ZeroCopyManager, ZeroCopyStats,
250};
251// `MemoryMappedBuffer` is mmap-backed and therefore `#[cfg(unix)]` in
252// `zero_copy`; the re-export has to carry the same gate or it fails to resolve
253// on Windows.
254#[cfg(unix)]
255pub use zero_copy::MemoryMappedBuffer;
256
257// New v0.3.0 exports for developer experience and performance
258pub 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
288// New v0.3.0 exports for ML, versioning, and migration
289pub 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
310// New v0.3.0 advanced ML exports
311pub 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
331// Utility exports
332pub use utils::{
333    create_dev_stream, create_prod_stream, BatchProcessor, EventFilter, EventSampler,
334    SimpleRateLimiter, StreamMultiplexer, StreamStats,
335};
336
337// Advanced SciRS2 optimization exports
338pub 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
346// Adaptive load shedding exports
347pub use adaptive_load_shedding::{
348    DropStrategy, LoadMetrics, LoadSheddingConfig, LoadSheddingManager, LoadSheddingStats,
349};
350
351// Stream fusion optimizer exports
352pub use stream_fusion::{
353    FusableChain, FusedOperation, FusedType, FusionAnalysis, FusionConfig, FusionOptimizer,
354    FusionStats, Operation,
355};
356
357// Complex Event Processing (CEP) engine exports
358pub 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
367// Data quality and validation framework exports
368pub 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
379// Advanced sampling techniques exports
380pub 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;
393/// REST client for an external Confluent-compatible Schema Registry. Pure Rust
394/// (reqwest/rustls) and independent of any broker — it outlived the removed
395/// rdkafka-backed Kafka backend for that reason.
396pub 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;
451// Standalone WASM plugin processor (distinct data model from `wasm_edge_computing`).
452// Its `execute_wasm_function` does not embed a real execution engine, so it always
453// returns `StreamError::UnsupportedOperation` rather than fabricating output; use
454// `wasm_edge_computing::WasmEdgeProcessor` (built on wasmtime) for genuine execution.
455pub mod wasm_edge_processor;
456pub mod webhook;
457
458// New v0.3.0 modules for advanced features
459pub mod custom_serialization;
460pub mod end_to_end_encryption;
461// PRE-EXISTING BUG WORKAROUND (unrelated to the rdkafka/pulsar quarantine):
462// `gpu_acceleration` imports `scirs2_core::gpu`, which scirs2-core gates behind its own
463// `gpu` feature (Pure Rust — `gpu = ["std"]`, no GPU FFI). Gate the module behind an
464// off-by-default `gpu` feature that turns scirs2-core's `gpu` on, so the default build
465// compiles without it and `--all-features` resolves `scirs2_core::gpu`.
466#[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
476// New v0.3.0 modules for developer experience and performance
477pub 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
486// New v0.3.0 modules for ML, versioning, and migration
487pub mod anomaly_detection;
488pub mod migration_tools;
489pub mod online_learning;
490pub mod stream_versioning;
491
492// Advanced ML modules for v0.3.0 completion
493pub mod automl_stream;
494pub mod feature_engineering;
495pub mod neural_architecture_search;
496pub mod predictive_analytics;
497pub mod reinforcement_learning;
498
499// Utilities module
500pub mod utils;
501
502// Advanced SciRS2 optimization module
503pub mod advanced_scirs2_optimization;
504pub mod cdc_processor;
505
506// Adaptive load shedding module
507pub mod adaptive_load_shedding;
508
509// Stream fusion optimizer module
510pub mod stream_fusion;
511
512// Complex Event Processing (CEP) engine module
513pub mod cep_engine;
514
515// Data quality and validation framework module
516pub mod data_quality;
517
518// Advanced sampling techniques module
519pub mod advanced_sampling;
520
521// Extracted type definitions to comply with 2000-line policy
522mod lib_types;
523pub use lib_types::*;
524
525// Distributed stream state management and fault tolerance (v0.2.0)
526pub mod checkpoint;
527
528// Re-export checkpoint types
529pub use checkpoint::{CheckpointCoordinator, CheckpointPhase, GlobalCheckpoint, OperatorSnapshot};
530
531// Re-export distributed state types
532pub 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
539// v0.3.0 modules: Distributed stream processing, distributed state, fault tolerance
540pub mod distributed;
541pub mod distributed_state;
542pub mod fault_tolerance;
543
544// WebSub (W3C) push notifications for RDF dataset changes
545pub mod websub;
546pub use websub::{
547    DatasetChangeEvent, DatasetEventBus, WebSubHub, WebSubPublisher, WebSubSubscriber,
548};
549
550// v1.0.0 ML module for streaming inference and anomaly detection
551pub mod ml;
552
553// Re-export ML types
554pub 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
561// Re-export distributed state manager types (v1.0.0)
562pub 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
570// Re-export fault tolerance checkpoint/recovery types (v1.0.0)
571pub use fault_tolerance::checkpoint_recovery::{
572    CheckpointManager, CheckpointManagerConfig, CheckpointManagerStats, CheckpointMetadata,
573    PartitionRebalancer, RebalanceAction, RebalancerConfig, RebalancerStats, RecoveryManager,
574    RecoveryManagerStats, RecoveryResult, RecoveryStrategy, StoredCheckpoint, WorkPartition,
575};
576
577// v0.2.0 modules: Distributed state, consistency protocols, watermarking, stream metrics
578pub mod consistency;
579pub mod metrics;
580pub mod watermark;
581
582// Re-export consistency types (v0.2.0)
583pub use consistency::ConsistencyManager as StreamConsistencyManager;
584pub use consistency::{
585    ConsistencyConfig, EventualConsistencyBuffer, StreamConsistencyLevel, VersionedValue,
586};
587
588// Re-export watermark types (v0.2.0)
589pub use watermark::{
590    LateDataDecision, LateDataHandler, LateDataPolicy, StreamWatermark, WatermarkAligner,
591    WatermarkGenerator,
592};
593
594// Re-export stream metrics types (v0.2.0)
595pub use metrics::{StreamLatencyHistogram, StreamMetrics, StreamMetricsCollector};
596
597// v1.1.0: Idempotent delivery with idempotency keys
598pub mod idempotent_delivery;
599
600// v1.1.0: Stream Windowing Algebra (tumbling, sliding, session, count-based)
601pub mod window_algebra;
602
603// v1.1.0 round 5: Adaptive backpressure controller (Drop/Block/Throttle/SpillToDisk)
604pub mod backpressure_controller;
605
606// v1.1.0 round 7: Synchronous in-memory schema registry (Kafka Schema Registry style)
607pub mod sync_schema_registry;
608
609// v1.1.0 round 11: Tumbling / sliding / session windowing functions
610pub mod window_function;
611
612// v1.1.0: Event sourcing patterns (append-only log, snapshotting, pub/sub bus)
613// Note: event_sourcing module already declared at line 401
614pub use idempotent_delivery::{
615    DeliveryOutcome, HashAlgorithm, IdempotencyKey, IdempotentDeliveryConfig,
616    IdempotentDeliveryManager, IdempotentDeliveryStats, IdempotentProducer, KeyCheckResult,
617};
618
619// v1.1.0 round 12: Stream checkpoint/offset tracking for at-least-once delivery
620pub mod stream_checkpoint;
621pub use stream_checkpoint::{Checkpoint, CheckpointStore};
622
623// v1.1.0 round 13: Dead letter queue for failed/undeliverable messages
624pub mod dead_letter_queue;
625
626// v1.1.0 round 14: Kafka-style consumer group coordination
627pub mod consumer_group;
628
629// v1.1.0 round 15: Stream message routing (content/topic/header/round-robin/DLQ)
630pub mod stream_router;
631
632// v1.1.0 round 16: Stream message schema validation (field types, formats, strict mode)
633pub mod schema_validator;
634
635// v1.1.0 round 17 (Batch E): Message format transformation pipeline
636pub mod message_transformer;
637
638// v1.1.0 round 18 (Batch E): Event replay buffer with seek and position tracking
639pub mod replay_buffer;
640
641// v1.1.0 round 19: Stream event filtering with composable predicates
642pub mod event_filter;
643
644// Visual pipeline designer and debugger (SVG/JSON/YAML/DOT/Mermaid export)
645pub mod visual_designer;
646pub mod visual_designer_engine;
647pub mod visual_designer_tests;
648pub mod visual_designer_types;
649
650// W2-S6: per-stream SLA admission control + load-shedder coordination.
651pub mod sla;
652pub use sla::{
653    BackpressureAction, SlaBackpressureCoordinator, SlaBackpressureDecision, SlaBackpressurePolicy,
654    StreamAdmissionController, StreamAdmissionDecision, StreamAdmissionStats, StreamSlaConfig,
655};
656
657// W2-S6: watermark-aware window joins (tumbling-tumbling, tumbling-sliding,
658// session-session) — a separate `window` module that complements the
659// time-based `processing::window`.
660pub mod window;
661pub use window::{
662    SessionSessionJoin, SessionSessionJoinConfig, TumblingSlidingJoin, TumblingSlidingJoinConfig,
663    TumblingTumblingJoin, TumblingTumblingJoinConfig, WindowJoinKey, WindowJoinResult,
664    WindowJoinStats,
665};
666
667// W2-S6: exactly-once aggregation under operator parallelism.
668pub mod aggregation;
669pub use aggregation::{
670    ExactlyOnceAggregator, ExactlyOnceAggregatorConfig, ExactlyOnceAggregatorStats,
671    PartitionAggregateState, PartitionAggregateValue,
672};
673
674// Neuromorphic stream analytics (brain-inspired spiking neural networks)
675pub 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;