Add eventing support to WMA indicator and implement unit tests for various indicators

- Enhanced WMA indicator with event-driven capabilities using ITValuePublisher interface.
- Created a new TODO file listing various indicators and their corresponding libraries.
- Added unit tests for DEMA, HMA, TEMA, and WMA indicators to ensure proper functionality.
- Implemented tests for handling new bars, ticks, and historical data updates across indicators.
- Verified that indicators correctly compute values and handle different source types.
This commit is contained in:
Miha Kralj
2025-12-07 16:46:38 -08:00
parent 3734a1c5f6
commit 875998b288
31 changed files with 2445 additions and 1457 deletions
+229
View File
@@ -0,0 +1,229 @@
using System;
using Xunit;
namespace QuanTAlib.Tests;
public class SmaCoverageTests
{
[Fact]
public void Sma_ResyncLogic_IsTriggeredAndCorrect()
{
// ResyncInterval is 1000.
int count = 2500;
int period = 10;
var sma = new Sma(period);
double constantValue = 100.0;
for (int i = 0; i < count; i++)
{
sma.Update(new TValue(DateTime.UtcNow, constantValue));
if (i >= period)
{
Assert.Equal(constantValue, sma.Last.Value, 1e-9);
}
}
}
[Fact]
public void Sma_SpanCalc_LargeDataset_TriggersResync()
{
int count = 5000;
int period = 10;
double[] source = new double[count];
double[] output = new double[count];
for (int i = 0; i < count; i++)
{
source[i] = 100.0;
}
Sma.Calculate(source.AsSpan(), output.AsSpan(), period);
for (int i = period; i < count; i++)
{
Assert.Equal(100.0, output[i], 1e-9);
}
}
[Fact]
public void Sma_SpanCalc_SimdThreshold_Boundary()
{
// SimdThreshold is 256.
int[] lengths = { 250, 256, 260 };
int period = 10;
foreach (int len in lengths)
{
double[] source = new double[len];
double[] output = new double[len];
for (int i = 0; i < len; i++) source[i] = 100.0;
Sma.Calculate(source.AsSpan(), output.AsSpan(), period);
Assert.Equal(100.0, output[^1], 1e-9);
}
}
[Fact]
public void Sma_SpanCalc_Simd_WithResync()
{
int count = 3000;
int period = 5;
double[] source = new double[count];
double[] output = new double[count];
// Linear increase: 0, 1, 2, ...
for (int i = 0; i < count; i++) source[i] = i;
Sma.Calculate(source.AsSpan(), output.AsSpan(), period);
// SMA(5) of x-4, x-3, x-2, x-1, x
// = (5x - 10) / 5 = x - 2
for (int i = period; i < count; i++)
{
double expected = i - 2.0;
Assert.Equal(expected, output[i], 1e-9);
}
}
[Fact]
public void Sma_Constructor_ThrowsOnInvalidPeriod()
{
Assert.Throws<ArgumentException>(() => new Sma(0));
Assert.Throws<ArgumentException>(() => new Sma(-1));
}
[Fact]
public void Sma_StaticCalculate_ThrowsOnInvalidArgs()
{
double[] source = new double[10];
double[] output = new double[5]; // Mismatch
Assert.Throws<ArgumentException>(() => Sma.Calculate(source.AsSpan(), output.AsSpan(), 5));
double[] output2 = new double[10];
Assert.Throws<ArgumentException>(() => Sma.Calculate(source.AsSpan(), output2.AsSpan(), 0));
}
[Fact]
public void Sma_Calculate_EmptyInput_DoesNothing()
{
Sma.Calculate(ReadOnlySpan<double>.Empty, Span<double>.Empty, 5);
// Should not throw
}
[Fact]
public void Sma_Update_WithNaN_UsesLastValid()
{
var sma = new Sma(5);
sma.Update(new TValue(DateTime.UtcNow, 1.0));
sma.Update(new TValue(DateTime.UtcNow, 2.0));
sma.Update(new TValue(DateTime.UtcNow, double.NaN)); // Should use 2.0
// Buffer: 1, 2, 2
// SMA(3) = (1 + 2 + 2) / 3 = 5/3 = 1.666...
Assert.Equal(5.0/3.0, sma.Last.Value, 1e-9);
}
[Fact]
public void Sma_Update_IsNewFalse_UpdatesLastValue()
{
var sma = new Sma(3);
sma.Update(new TValue(DateTime.UtcNow, 1.0));
sma.Update(new TValue(DateTime.UtcNow, 2.0));
// Update existing with 3.0 (replaces 2.0)
sma.Update(new TValue(DateTime.UtcNow, 3.0), isNew: false);
// Buffer should be: 1, 3
// SMA = (1 + 3) / 2 = 2
Assert.Equal(2.0, sma.Last.Value, 1e-9);
}
[Fact]
public void Sma_TSeries_Empty_ReturnsEmpty()
{
var sma = new Sma(5);
var result = sma.Update(new TSeries());
Assert.Empty(result);
}
[Fact]
public void Sma_TSeries_WithNaN_RestoresStateCorrectly()
{
var sma = new Sma(3);
var series = new TSeries();
series.Add(new TValue(DateTime.UtcNow, 1.0));
series.Add(new TValue(DateTime.UtcNow, 2.0));
series.Add(new TValue(DateTime.UtcNow, double.NaN));
series.Add(new TValue(DateTime.UtcNow, 4.0));
sma.Update(series);
// Buffer: 2.0, 2.0 (from NaN), 4.0
// SMA(3) = (2 + 2 + 4) / 3 = 8/3 = 2.666...
// Let's add one more value to verify state is correct
sma.Update(new TValue(DateTime.UtcNow, 5.0));
// Buffer: 2.0, 4.0, 5.0
// SMA(3) = (2 + 4 + 5) / 3 = 11/3 = 3.666...
Assert.Equal(11.0/3.0, sma.Last.Value, 1e-9);
}
[Fact]
public void Sma_Reset_ClearsState()
{
var sma = new Sma(3);
sma.Update(new TValue(DateTime.UtcNow, 1.0));
sma.Update(new TValue(DateTime.UtcNow, 2.0));
sma.Update(new TValue(DateTime.UtcNow, 3.0));
sma.Reset();
Assert.Equal(0, sma.Last.Value);
// Start fresh
sma.Update(new TValue(DateTime.UtcNow, 10.0));
// Buffer: 10
// SMA = 10
Assert.Equal(10.0, sma.Last.Value);
}
[Fact]
public void Sma_Calculate_ScalarFallback_WithNaN()
{
// Force scalar path by including NaN, even with large dataset
int count = 1000;
double[] source = new double[count];
double[] output = new double[count];
for (int i = 0; i < count; i++) source[i] = 1.0;
source[500] = double.NaN; // This should trigger HasNonFiniteValues -> true
Sma.Calculate(source.AsSpan(), output.AsSpan(), 10);
// Check around the NaN
// Index 500 is NaN, so it uses previous valid (1.0)
// So effectively the stream is all 1.0s
Assert.Equal(1.0, output[500], 1e-9);
Assert.Equal(1.0, output[501], 1e-9);
}
[Fact]
public void Sma_Constructor_WithSource_Subscribes()
{
var source = new Sma(10); // Just using Sma as a publisher
var sma = new Sma(source, 5);
source.Update(new TValue(DateTime.UtcNow, 10.0));
Assert.Equal(10.0, sma.Last.Value);
}
}
-285
View File
@@ -1,285 +0,0 @@
#!meta
{"kernelInfo":{"defaultKernelName":"csharp","items":[{"name":"csharp"},{"name":"fsharp","languageName":"F#","aliases":["f#","fs"]},{"name":"html","languageName":"HTML"},{"name":"http","languageName":"HTTP"},{"name":"javascript","languageName":"JavaScript","aliases":["js"]},{"name":"mermaid","languageName":"Mermaid"},{"name":"pwsh","languageName":"PowerShell","aliases":["powershell"]},{"name":"value"}]}}
#!markdown
# Simple Moving Average (SMA) Examples
This is a **.NET Interactive** notebook. To run it, you need the [Polyglot Notebooks](https://marketplace.visualstudio.com/items?itemName=ms-dotnettools.dotnet-interactive-vscode) extension installed in VS Code.
The **Simple Moving Average (SMA)** is the most basic form of moving average, calculating the arithmetic mean over a specified period. Unlike the EMA, the SMA assigns equal weight to all data points in the window, making it a good baseline for trend analysis.
**Key characteristics:**
- Equal weighting for all values in the period
- O(1) update complexity using running sum
- O(1) bar correction using scalar state
- Smooth output with good noise reduction
- More lag than EMA due to equal weighting
This notebook demonstrates:
1. **Manual Data Processing**: Understanding Batch vs. Streaming modes.
2. **Streaming with `isNew`**: Handling intra-bar updates.
3. **Large Dataset Processing**: Using Geometric Brownian Motion (GBM) generated data.
4. **Handling Invalid Values**: Last-value substitution for NaN/Infinity.
5. **SMA vs EMA**: Comparing Simple and Exponential Moving Averages.
#!csharp
// Reference the library
#r "..\..\bin\QuanTAlib.dll"
using System;
using System.Linq;
using QuanTAlib;
// Helper to print TSeries
void PrintSeries(TSeries series, int count = 5)
{
Console.WriteLine($"Series Length: {series.Count}");
foreach (var item in series.Take(count))
{
Console.WriteLine($"Time: {item.Time:HH:mm:ss}, Value: {item.Value:F2}");
}
if (series.Count > count) Console.WriteLine("...");
}
#!markdown
## 1. Manual Data: Batch vs. Streaming
We'll start with a small, manually created dataset to clearly see how Batch and Streaming operations work.
### Batch Processing
Batch processing calculates the SMA for the entire dataset at once. This is efficient for historical analysis.
#!csharp
// Create a small manual dataset
var manualData = new TSeries();
manualData.Add(DateTime.Now, 100.0);
manualData.Add(DateTime.Now.AddMinutes(1), 102.0);
manualData.Add(DateTime.Now.AddMinutes(2), 101.0);
manualData.Add(DateTime.Now.AddMinutes(3), 103.0);
manualData.Add(DateTime.Now.AddMinutes(4), 105.0);
Console.WriteLine("--- Input Data ---");
PrintSeries(manualData, 5);
// Batch Calculation
Console.WriteLine("\n--- Batch SMA (Period 3) ---");
var smaBatch = new Sma(3);
var resultBatch = smaBatch.Update(manualData);
PrintSeries(resultBatch, 5);
// Show the calculation for each step
Console.WriteLine("\nCalculation breakdown:");
Console.WriteLine(" SMA[0] = 100 / 1 = 100.00");
Console.WriteLine(" SMA[1] = (100 + 102) / 2 = 101.00");
Console.WriteLine(" SMA[2] = (100 + 102 + 101) / 3 = 101.00");
Console.WriteLine(" SMA[3] = (102 + 101 + 103) / 3 = 102.00");
Console.WriteLine(" SMA[4] = (101 + 103 + 105) / 3 = 103.00");
#!markdown
### Streaming Processing
Streaming processing updates the SMA one data point at a time. This is essential for real-time trading systems where data arrives sequentially.
#!csharp
Console.WriteLine("\n--- Streaming SMA (Period 3) ---");
var smaStream = new Sma(3);
foreach (var item in manualData)
{
var result = smaStream.Update(item);
Console.WriteLine($"Time: {item.Time:HH:mm:ss}, Input: {item.Value:F2}, SMA: {result.Value:F2}, IsHot: {smaStream.IsHot}");
}
// Verify that the last values match
var batchLast = resultBatch.Last().Value;
var streamLast = smaStream.Value.Value;
Console.WriteLine($"\nMatch: {Math.Abs(batchLast - streamLast) < 1e-10} (Batch: {batchLast:F2}, Stream: {streamLast:F2})");
// Show SMA properties
Console.WriteLine($"\nSMA Properties:");
Console.WriteLine($" Name: {smaStream.Name}");
Console.WriteLine($" WarmupPeriod: {smaStream.WarmupPeriod}");
Console.WriteLine($" IsHot: {smaStream.IsHot}");
#!markdown
## 2. Streaming with `isNew` (Intra-bar Updates)
In real-time feeds, you often receive multiple updates for the *same* bar (e.g., price changes within the current minute) before the bar closes.
* `isNew = true`: The input is a new bar (advances time).
* `isNew = false`: The input is an update to the current bar (recalculates without advancing).
**SMA achieves O(1) bar correction** by saving scalar state after each `isNew=true` update.
#!csharp
Console.WriteLine("\n--- Streaming with Intra-bar Updates ---");
var smaIntra = new Sma(3);
// 1. Process the first 4 bars normally
for (int i = 0; i < 4; i++)
{
smaIntra.Update(manualData[i]);
}
Console.WriteLine($"After 4th bar: {smaIntra.Value.Value:F2}");
// 2. Simulate intra-bar updates for the 5th bar (Final value is 105.0)
// Update 1: Price moves to 104.0
var update1 = new TValue(manualData[4].Time, 104.0);
smaIntra.Update(update1, isNew: true); // First update for this bar is "New"
Console.WriteLine($"Update 1 (104.0): {smaIntra.Value.Value:F2}");
// Update 2: Price moves to 106.0 (Same time, same bar)
var update2 = new TValue(manualData[4].Time, 106.0);
smaIntra.Update(update2, isNew: false); // Not new, just an update
Console.WriteLine($"Update 2 (106.0): {smaIntra.Value.Value:F2}");
// Update 3: Final Close at 105.0
var update3 = manualData[4];
smaIntra.Update(update3, isNew: false); // Final update
Console.WriteLine($"Update 3 (105.0): {smaIntra.Value.Value:F2}");
// Verify match with batch result
Console.WriteLine($"Match with Batch: {Math.Abs(smaIntra.Value.Value - batchLast) < 1e-10}");
#!markdown
## 3. Large Dataset: Geometric Brownian Motion (GBM)
We'll generate a larger dataset (1000 bars) using a Geometric Brownian Motion generator to simulate realistic market data.
#!csharp
// Generate 1000 bars of data
var gbm = new GBM(startPrice: 100.0, mu: 0.05, sigma: 0.2);
var gbmData = gbm.Fetch(1000, DateTime.Now.Ticks, TimeSpan.FromMinutes(1));
var closeSeries = gbmData.Close;
Console.WriteLine($"Generated {closeSeries.Count} bars of GBM data.");
Console.WriteLine($"First 5 values: {string.Join(", ", closeSeries.Take(5).Select(x => x.Value.ToString("F2")))}");
#!markdown
### Batch vs. Streaming Performance on Large Data
#!csharp
// Batch
var smaLargeBatch = new Sma(20);
var batchLargeResult = smaLargeBatch.Update(closeSeries);
Console.WriteLine($"Batch Last Value: {batchLargeResult.Last().Value:F2}");
// Streaming
var smaLargeStream = new Sma(20);
TValue lastStreamVal = default;
foreach(var item in closeSeries)
{
lastStreamVal = smaLargeStream.Update(item);
}
Console.WriteLine($"Streaming Last Value: {lastStreamVal.Value:F2}");
// Verify match
Console.WriteLine($"Match: {Math.Abs(batchLargeResult.Last().Value - lastStreamVal.Value) < 1e-10}");
#!markdown
## 4. Handling Invalid Values (NaN/Infinity)
`Sma` uses **last-value substitution** for invalid inputs. When a non-finite value (NaN, PositiveInfinity, NegativeInfinity) is encountered, it is replaced with the last valid value. This provides output continuity instead of propagating invalid values through the calculation.
#!csharp
Console.WriteLine("\n--- Handling Invalid Values ---");
// Single SMA
var smaNaN = new Sma(10);
// Feed valid values first
smaNaN.Update(new TValue(DateTime.Now, 100.0));
smaNaN.Update(new TValue(DateTime.Now.AddMinutes(1), 110.0));
Console.WriteLine($"After valid values: {smaNaN.Value.Value:F2}");
// Feed NaN - should use last valid value (110)
var resultAfterNaN = smaNaN.Update(new TValue(DateTime.Now.AddMinutes(2), double.NaN));
Console.WriteLine($"After NaN input: {resultAfterNaN.Value:F2} (IsFinite: {double.IsFinite(resultAfterNaN.Value)})");
// Feed Infinity - should use last valid value (110)
var resultAfterInf = smaNaN.Update(new TValue(DateTime.Now.AddMinutes(3), double.PositiveInfinity));
Console.WriteLine($"After Infinity input: {resultAfterInf.Value:F2} (IsFinite: {double.IsFinite(resultAfterInf.Value)})");
// Continue with valid value
var resultAfterValid = smaNaN.Update(new TValue(DateTime.Now.AddMinutes(4), 120.0));
Console.WriteLine($"After valid value (120): {resultAfterValid.Value:F2}");
#!csharp
Console.WriteLine("\n--- Batch Processing with Invalid Values ---");
// Create series with NaN values interspersed
var seriesWithNaN = new TSeries();
seriesWithNaN.Add(DateTime.Now.Ticks, 100.0);
seriesWithNaN.Add(DateTime.Now.Ticks + 1, 110.0);
seriesWithNaN.Add(DateTime.Now.Ticks + 2, double.NaN);
seriesWithNaN.Add(DateTime.Now.Ticks + 3, 120.0);
seriesWithNaN.Add(DateTime.Now.Ticks + 4, double.PositiveInfinity);
seriesWithNaN.Add(DateTime.Now.Ticks + 5, 130.0);
var smaBatchNaN = new Sma(3);
var resultsWithNaN = smaBatchNaN.Update(seriesWithNaN);
Console.WriteLine("Input → Output:");
for (int i = 0; i < seriesWithNaN.Count; i++)
{
var input = seriesWithNaN[i].Value;
var output = resultsWithNaN[i].Value;
var inputStr = double.IsFinite(input) ? input.ToString("F2") : input.ToString();
Console.WriteLine($" {inputStr,-10} → {output:F2} (IsFinite: {double.IsFinite(output)})");
}
#!markdown
## 5. SMA vs EMA Comparison
The SMA and EMA are both trend-following indicators, but they weight data differently:
- **SMA**: Equal weight to all values in the window
- **EMA**: More weight to recent values (exponentially decreasing)
#!csharp
Console.WriteLine("\n--- SMA vs EMA Comparison (Period 10) ---");
var compareData = new TSeries();
var baseTime = DateTime.Now;
for (int i = 0; i < 20; i++)
{
// Create data with a sudden spike at position 10
double value = (i == 10) ? 150.0 : 100.0;
compareData.Add(baseTime.AddMinutes(i), value);
}
var smaCompare = new Sma(10);
var emaCompare = new Ema(10);
Console.WriteLine("Position | Input | SMA | EMA | Difference");
Console.WriteLine("---------+--------+---------+---------+-----------");
for (int i = 0; i < compareData.Count; i++)
{
var smaVal = smaCompare.Update(compareData[i]);
var emaVal = emaCompare.Update(compareData[i]);
var input = compareData[i].Value;
var diff = smaVal.Value - emaVal.Value;
Console.WriteLine($" {i,2} | {input,6:F0} | {smaVal.Value,7:F2} | {emaVal.Value,7:F2} | {diff,+9:F2}");
}
Console.WriteLine("\nNote: After the spike (position 10), EMA reacts faster due to higher weight on recent values.");
Console.WriteLine("SMA takes longer to reflect changes as all values have equal weight.");
+35 -30
View File
@@ -108,25 +108,13 @@ public sealed class Sma : ITValuePublisher
}
// Removed GetValidValue and UpdateState as they are not used in the new Update logic
[MethodImpl(MethodImplOptions.AggressiveInlining)]
public TValue Update(TValue input, bool isNew = true)
{
if (isNew)
{
double val = GetValidValue(input.Value);
double removedValue = _buffer.Count == _buffer.Capacity ? _buffer.Oldest : 0.0;
_sum = _sum - removedValue + val;
_buffer.Add(val);
_tickCount++;
if (_buffer.IsFull && _tickCount >= ResyncInterval)
{
_tickCount = 0;
_sum = _buffer.Sum();
}
UpdateState(val);
_p_sum = _sum;
_p_lastInput = val;
@@ -136,7 +124,7 @@ public sealed class Sma : ITValuePublisher
{
_lastValidValue = _p_lastValidValue;
double val = GetValidValue(input.Value);
_sum = _p_sum - _p_lastInput + val;
_buffer.UpdateNewest(val);
}
@@ -150,7 +138,7 @@ public sealed class Sma : ITValuePublisher
public TSeries Update(TSeries source)
{
if (source.Count == 0) return new TSeries(new List<long>(), new List<double>());
int len = source.Count;
var t = new List<long>(len);
var v = new List<double>(len);
@@ -159,26 +147,43 @@ public sealed class Sma : ITValuePublisher
var tSpan = CollectionsMarshal.AsSpan(t);
var vSpan = CollectionsMarshal.AsSpan(v);
var sourceValues = source.Values;
var sourceTimes = source.Times;
// Reset state for batch calculation
Reset();
// We can optimize this later with specific batch logic, but for now use core loop
for(int i=0; i < len; i++)
Calculate(source.Values, vSpan, _period);
source.Times.CopyTo(tSpan);
// Restore state
int windowSize = Math.Min(len, _period);
int startIndex = len - windowSize;
if (startIndex > 0)
{
double val = GetValidValue(sourceValues[i]);
double removedValue = _buffer.Count == _buffer.Capacity ? _buffer.Oldest : 0.0;
_sum = _sum - removedValue + val;
_buffer.Add(val);
vSpan[i] = _sum / _buffer.Count;
for (int i = startIndex - 1; i >= 0; i--)
{
if (double.IsFinite(source.Values[i]))
{
_lastValidValue = source.Values[i];
break;
}
}
}
else
{
_lastValidValue = 0;
}
_buffer.Clear();
_sum = 0;
_tickCount = 0;
for (int i = startIndex; i < len; i++)
{
double val = GetValidValue(source.Values[i]);
UpdateState(val);
}
sourceTimes.CopyTo(tSpan);
_p_lastValidValue = _lastValidValue;
_p_sum = _sum;
_p_lastInput = sourceValues[len-1];
_p_lastInput = source.Values[len - 1];
_p_lastValidValue = _lastValidValue;
Last = new TValue(tSpan[len - 1], vSpan[len - 1]);
return new TSeries(t, v);
+28
View File
@@ -125,6 +125,34 @@ sma.Update(new TValue(time + 1, 101.2), isNew: true);
**Implementation detail:** Bar correction is O(1) using scalar state save/restore, not buffer copying.
### Eventing and Reactive Support
This indicator implements the `ITValuePublisher` interface, enabling event-driven and reactive workflows.
* **Subscription:** Can be constructed with an `ITValuePublisher` (e.g., `TSeries`) to automatically update when the source emits a new value.
* **Publication:** Emits a `Pub` event with the new `TValue` whenever it is updated.
```csharp
using QuanTAlib;
// 1. Setup a source (publisher)
var source = new TSeries();
// 2. Create indicator subscribed to source
// It waits for events from 'source'
var sma = new Sma(source, period: 10);
// 3. Optional: Subscribe to indicator's output
sma.Pub += (item) => Console.WriteLine($"SMA Updated: {item.Value}");
// 4. Ingest data into source
// This triggers the chain: source -> sma -> Console.WriteLine
source.Add(new TValue(DateTime.Now, 100));
source.Add(new TValue(DateTime.Now, 105));
```
This pattern allows building complex, reactive processing pipelines without manual update loops.
### Handling Invalid Values (NaN/Infinity)
`Sma` uses **last-value substitution** for handling invalid inputs: