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.3-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.3)
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;
442// Standalone WASM plugin processor (distinct data model from `wasm_edge_computing`).
443// Its `execute_wasm_function` does not embed a real execution engine, so it always
444// returns `StreamError::UnsupportedOperation` rather than fabricating output; use
445// `wasm_edge_computing::WasmEdgeProcessor` (built on wasmtime) for genuine execution.
446pub mod wasm_edge_processor;
447pub mod webhook;
448
449// New v0.3.0 modules for advanced features
450pub mod custom_serialization;
451pub mod end_to_end_encryption;
452// PRE-EXISTING BUG WORKAROUND (unrelated to the rdkafka/pulsar quarantine):
453// `gpu_acceleration` imports `scirs2_core::gpu`, which scirs2-core gates behind its own
454// `gpu` feature (Pure Rust — `gpu = ["std"]`, no GPU FFI). Gate the module behind an
455// off-by-default `gpu` feature that turns scirs2-core's `gpu` on, so the default build
456// compiles without it and `--all-features` resolves `scirs2_core::gpu`.
457#[cfg(feature = "gpu")]
458pub mod gpu_acceleration;
459pub mod ml_integration;
460pub mod rate_limiting;
461pub mod scalability;
462pub mod schema_evolution;
463pub mod stream_replay;
464pub mod transactional_processing;
465pub mod zero_copy;
466
467// New v0.3.0 modules for developer experience and performance
468pub mod numa_processing;
469pub mod out_of_order;
470pub mod performance_profiler;
471pub mod stream_sql;
472pub mod stream_sql_ast;
473pub mod stream_sql_executor;
474pub mod stream_sql_tests;
475pub mod testing_framework;
476
477// New v0.3.0 modules for ML, versioning, and migration
478pub mod anomaly_detection;
479pub mod migration_tools;
480pub mod online_learning;
481pub mod stream_versioning;
482
483// Advanced ML modules for v0.3.0 completion
484pub mod automl_stream;
485pub mod feature_engineering;
486pub mod neural_architecture_search;
487pub mod predictive_analytics;
488pub mod reinforcement_learning;
489
490// Utilities module
491pub mod utils;
492
493// Advanced SciRS2 optimization module
494pub mod advanced_scirs2_optimization;
495pub mod cdc_processor;
496
497// Adaptive load shedding module
498pub mod adaptive_load_shedding;
499
500// Stream fusion optimizer module
501pub mod stream_fusion;
502
503// Complex Event Processing (CEP) engine module
504pub mod cep_engine;
505
506// Data quality and validation framework module
507pub mod data_quality;
508
509// Advanced sampling techniques module
510pub mod advanced_sampling;
511
512// Extracted type definitions to comply with 2000-line policy
513mod lib_types;
514pub use lib_types::*;
515
516// Distributed stream state management and fault tolerance (v0.2.0)
517pub mod checkpoint;
518
519// Re-export checkpoint types
520pub use checkpoint::{CheckpointCoordinator, CheckpointPhase, GlobalCheckpoint, OperatorSnapshot};
521
522// Re-export distributed state types
523pub use state::{
524    AggregatingState, DeduplicationConfig, DeduplicationLog, DistributedStateBackend,
525    DistributedStateStore, ExactlyOnceProcessor, ExactlyOnceStats, ExactlyOnceTransaction,
526    InMemoryStateBackend, KeyedStateStore, MessageId, PartitionStateValue, StateAggregator,
527    StateBackendStats, StateCoordinator, StatePartition, StatePartitionKey,
528};
529
530// v0.3.0 modules: Distributed stream processing, distributed state, fault tolerance
531pub mod distributed;
532pub mod distributed_state;
533pub mod fault_tolerance;
534
535// WebSub (W3C) push notifications for RDF dataset changes
536pub mod websub;
537pub use websub::{
538    DatasetChangeEvent, DatasetEventBus, WebSubHub, WebSubPublisher, WebSubSubscriber,
539};
540
541// v1.0.0 ML module for streaming inference and anomaly detection
542pub mod ml;
543
544// Re-export ML types
545pub use ml::{
546    AnomalyCheckResult, AnomalyDetectorConfig, AnomalyDetectorStats, ExtractedFeatures,
547    FeatureAggregation, FeatureDefinition, FeatureExtractorConfig, ModelConfig, ModelRunnerStats,
548    Prediction as MlPrediction, StreamAnomalyDetector,
549    StreamFeatureExtractor as MlStreamFeatureExtractor, StreamingModelRunner,
550};
551
552// Re-export distributed state manager types (v1.0.0)
553pub use distributed_state::manager::{
554    CheckpointConfig as DistributedCheckpointConfig,
555    DeduplicationConfig as DistributedDeduplicationConfig, DeduplicationStats,
556    DistributedStateManager, DistributedStateManagerStats, MigrationPlan, MigrationReason,
557    MigrationStep, OperatorStateSnapshot, PartitionAssignment, SequenceDeduplicator,
558    StateCheckpoint,
559};
560
561// Re-export fault tolerance checkpoint/recovery types (v1.0.0)
562pub use fault_tolerance::checkpoint_recovery::{
563    CheckpointManager, CheckpointManagerConfig, CheckpointManagerStats, CheckpointMetadata,
564    PartitionRebalancer, RebalanceAction, RebalancerConfig, RebalancerStats, RecoveryManager,
565    RecoveryManagerStats, RecoveryResult, RecoveryStrategy, StoredCheckpoint, WorkPartition,
566};
567
568// v0.2.0 modules: Distributed state, consistency protocols, watermarking, stream metrics
569pub mod consistency;
570pub mod metrics;
571pub mod watermark;
572
573// Re-export consistency types (v0.2.0)
574pub use consistency::ConsistencyManager as StreamConsistencyManager;
575pub use consistency::{
576    ConsistencyConfig, EventualConsistencyBuffer, StreamConsistencyLevel, VersionedValue,
577};
578
579// Re-export watermark types (v0.2.0)
580pub use watermark::{
581    LateDataDecision, LateDataHandler, LateDataPolicy, StreamWatermark, WatermarkAligner,
582    WatermarkGenerator,
583};
584
585// Re-export stream metrics types (v0.2.0)
586pub use metrics::{StreamLatencyHistogram, StreamMetrics, StreamMetricsCollector};
587
588// v1.1.0: Idempotent delivery with idempotency keys
589pub mod idempotent_delivery;
590
591// v1.1.0: Stream Windowing Algebra (tumbling, sliding, session, count-based)
592pub mod window_algebra;
593
594// v1.1.0 round 5: Adaptive backpressure controller (Drop/Block/Throttle/SpillToDisk)
595pub mod backpressure_controller;
596
597// v1.1.0 round 7: Synchronous in-memory schema registry (Kafka Schema Registry style)
598pub mod sync_schema_registry;
599
600// v1.1.0 round 11: Tumbling / sliding / session windowing functions
601pub mod window_function;
602
603// v1.1.0: Event sourcing patterns (append-only log, snapshotting, pub/sub bus)
604// Note: event_sourcing module already declared at line 401
605pub use idempotent_delivery::{
606    DeliveryOutcome, HashAlgorithm, IdempotencyKey, IdempotentDeliveryConfig,
607    IdempotentDeliveryManager, IdempotentDeliveryStats, IdempotentProducer, KeyCheckResult,
608};
609
610// v1.1.0 round 12: Stream checkpoint/offset tracking for at-least-once delivery
611pub mod stream_checkpoint;
612pub use stream_checkpoint::{Checkpoint, CheckpointStore};
613
614// v1.1.0 round 13: Dead letter queue for failed/undeliverable messages
615pub mod dead_letter_queue;
616
617// v1.1.0 round 14: Kafka-style consumer group coordination
618pub mod consumer_group;
619
620// v1.1.0 round 15: Stream message routing (content/topic/header/round-robin/DLQ)
621pub mod stream_router;
622
623// v1.1.0 round 16: Stream message schema validation (field types, formats, strict mode)
624pub mod schema_validator;
625
626// v1.1.0 round 17 (Batch E): Message format transformation pipeline
627pub mod message_transformer;
628
629// v1.1.0 round 18 (Batch E): Event replay buffer with seek and position tracking
630pub mod replay_buffer;
631
632// v1.1.0 round 19: Stream event filtering with composable predicates
633pub mod event_filter;
634
635// Visual pipeline designer and debugger (SVG/JSON/YAML/DOT/Mermaid export)
636pub mod visual_designer;
637pub mod visual_designer_engine;
638pub mod visual_designer_tests;
639pub mod visual_designer_types;
640
641// W2-S6: per-stream SLA admission control + load-shedder coordination.
642pub mod sla;
643pub use sla::{
644    BackpressureAction, SlaBackpressureCoordinator, SlaBackpressureDecision, SlaBackpressurePolicy,
645    StreamAdmissionController, StreamAdmissionDecision, StreamAdmissionStats, StreamSlaConfig,
646};
647
648// W2-S6: watermark-aware window joins (tumbling-tumbling, tumbling-sliding,
649// session-session) — a separate `window` module that complements the
650// time-based `processing::window`.
651pub mod window;
652pub use window::{
653    SessionSessionJoin, SessionSessionJoinConfig, TumblingSlidingJoin, TumblingSlidingJoinConfig,
654    TumblingTumblingJoin, TumblingTumblingJoinConfig, WindowJoinKey, WindowJoinResult,
655    WindowJoinStats,
656};
657
658// W2-S6: exactly-once aggregation under operator parallelism.
659pub mod aggregation;
660pub use aggregation::{
661    ExactlyOnceAggregator, ExactlyOnceAggregatorConfig, ExactlyOnceAggregatorStats,
662    PartitionAggregateState, PartitionAggregateValue,
663};
664
665// Neuromorphic stream analytics (brain-inspired spiking neural networks)
666pub mod neuromorphic_analytics;
667pub mod neuromorphic_analytics_engine;
668pub mod neuromorphic_analytics_learning;
669pub mod neuromorphic_analytics_network;
670pub mod neuromorphic_analytics_patterns;
671mod neuromorphic_analytics_tests;
672pub mod neuromorphic_analytics_types;