diff --git a/.github/workflows/main.yml b/.github/workflows/main.yml index 3d482c4..ae8cbeb 100644 --- a/.github/workflows/main.yml +++ b/.github/workflows/main.yml @@ -19,3 +19,18 @@ jobs: run: python -m pip install platformio - name: Build firmware run: pio ci --lib="." --board ${{matrix.pio-env}} "examples/${{matrix.example}}" + + test: + name: Unit Tests + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v2 + - name: Set up python + uses: actions/setup-python@v2 + with: + python-version: '3.x' + architecture: 'x64' + - name: Install PlatformIO + run: python -m pip install platformio + - name: Run unit tests + run: pio test -e native diff --git a/README.md b/README.md index ab3591a..fecd5ef 100644 --- a/README.md +++ b/README.md @@ -23,7 +23,137 @@ More examples in Wiki/[Examples](https://github.com/luisllamasbinaburo/Arduino-R ### Observable, observers and operators legend More info about the Observables, Observers, and Operators available in the [Wiki](https://github.com/luisllamasbinaburo/Arduino-ReactiveArduino/wiki) -![alt text](https://github.com/luisllamasbinaburo/Arduino-ReactiveArduino/blob/master/ReactiveArduino%20Legend.png "Legend") +```mermaid +flowchart LR + subgraph OBS[Observables] + direction TB + ManualDefer + Range + RangeDefer + Array + ArrayDefer + Property + AnalogInput + DigitalInput + TimerMillis + TimerMicros + IntervalMillis + IntervalMicros + SerialChar + SerialByte + SerialString + SerialInteger + SerialFloat + SerialDouble + end + + subgraph OPR[Operators] + direction TB + subgraph OPO[Operators] + Where + Distinct + First + Last + Take + TakeAt + TakeFirst + TakeLast + TakeUntil + TakeWhile + Skip + SkipUntil + SkipWhile + Batch + TimeoutMillis + TimeoutMicros + ForEach + If + Loop + Repeat + Reset + NoReset + end + subgraph TRN[Transformations] + Select + Cast + Map + Reduce + Limit + LimitLower + LimitUpper + Scale + ElapsedMicros + ElapsedMillis + Micros + Millis + Frequency + Threshold + Toggle + AdcToVoltage + Split + Join + Buffer + StringBuffer + ToBool + ToInt + ToFloat + ParseInt + ParseFloat + end + subgraph FLT[Filters] + OnRising + OnFalling + Median3 + Median5 + MovingAverage + MovingRMS + LowPass + HighPass + PassBand + StopBand + WindowMillis + WindowMicros + DebounceMillis + DebounceMicros + IsLessOrEqual + IsLess + IsGreaterOrEqual + IsGreater + IsNotEqual + IsEqual + IsZero + IsNotZero + end + subgraph AGG[Aggregates] + Count + Countdown + Sum + Min + Max + Average + Any + RMS + All + None + end + end + + subgraph OBV[Observers] + direction TB + Do + Finally + DoAndFinally + DoNothing + Property + Array + CircularBuffer + DigitalOutput + AnalogOutput + Serial + end + + OBS --> OPR --> OBV +``` ### Creating Observables Observables are generally generated through factory methods provided by the Reactive class. @@ -51,10 +181,19 @@ Hot observables emits the sequence when an observer subscribes to it. For exampl ```c++ FromArray(values, valuesLength) ``` -Cold observable does not emits any item when a observer subscribes to it. You have to explicitly call the `Next()` method whenever you want. For example, `FromArrayDefer(...)` +Cold observable does not emits any item when a observer subscribes to it. You have to explicitly call the `Next()` method whenever you want. For example, `FromArrayDefer(...)` or `ManualDefer(...)`. ```c++ FromArrayDefer(values, valuesLength) ``` + +With `ManualDefer(...)` you also have to explicitly call `Next()` to emit, and `Complete()` to finish the sequence. +```c++ +auto obs = ManualDefer(); +obs >> ToSerial(); +// ...later in code +obs.Next(); +obs.Complete(); +``` ### Dynamic memory considerations On many occasions we generate operators directly when we chain them, for example in the `Setup()`. However, creating an operator allocates dynamic memory. Therefore, you should avoid creating them in `Loop()`, or you could run out of memory. If you need to reuse (typically, call some operator method later in your code) set it as a global variable, and chain as normal. diff --git a/ReactiveArduino Legend.png b/ReactiveArduino Legend.png deleted file mode 100644 index c6335cf..0000000 Binary files a/ReactiveArduino Legend.png and /dev/null differ diff --git a/platformio.ini b/platformio.ini new file mode 100644 index 0000000..979410b --- /dev/null +++ b/platformio.ini @@ -0,0 +1,19 @@ +; PlatformIO Unit Testing configuration for the ReactiveArduino library. +; +; Run the tests with: +; pio test -e native +; +; The `native` environment compiles and runs the tests on the host machine. +; Arduino hardware functions (millis, micros, pinMode, digitalWrite, analogWrite, +; analogRead, digitalRead, String, Serial) are stubbed in test/WProgram.h so the +; header-only library can be exercised without a board. + +[platformio] +default_envs = native + +[env:native] +platform = native +test_framework = unity +build_flags = + -Isrc + -Itest diff --git a/src/Aggregates/AggregateAll.h b/src/Aggregates/AggregateAll.h index c952ac0..1a6b2b6 100644 --- a/src/Aggregates/AggregateAll.h +++ b/src/Aggregates/AggregateAll.h @@ -33,7 +33,7 @@ AggregateAll::AggregateAll(ReactivePredicate condition) template void AggregateAll::OnNext(T value) { - if (_state && _condition(value)) _state = false; + if (!_condition(value)) _state = false; this->_childObservers.OnNext(_state); } diff --git a/src/Aggregates/AggregateAverage.h b/src/Aggregates/AggregateAverage.h index 47d11a5..d030f0c 100644 --- a/src/Aggregates/AggregateAverage.h +++ b/src/Aggregates/AggregateAverage.h @@ -33,7 +33,7 @@ void AggregateAverage::OnNext(T value) { _sum += value; _count++; - this->_childObservers.OnNext(_sum / _count); + this->_childObservers.OnNext(_sum / static_cast(_count)); } #endif \ No newline at end of file diff --git a/src/Aggregates/AggregateCount.h b/src/Aggregates/AggregateCount.h index 95b6d07..54ff497 100644 --- a/src/Aggregates/AggregateCount.h +++ b/src/Aggregates/AggregateCount.h @@ -19,7 +19,7 @@ class AggregateCount : public Operator void OnNext(T value) override; private: - int _count = false; + int _count = 0; }; template diff --git a/src/Aggregates/AggregateCountdown.h b/src/Aggregates/AggregateCountdown.h index edc4efc..e625727 100644 --- a/src/Aggregates/AggregateCountdown.h +++ b/src/Aggregates/AggregateCountdown.h @@ -19,7 +19,8 @@ class AggregateCountdown : public Operator void OnNext(T value) override; private: - int _count = false; + int _count = 0; + bool _completed = false; }; template @@ -31,11 +32,14 @@ AggregateCountdown::AggregateCountdown(int count) template void AggregateCountdown::OnNext(T value) { + if (_completed) return; + _count--; this->_childObservers.OnNext(_count); if (_count <= 0) { + _completed = true; this->_childObservers.OnComplete(); } } diff --git a/src/Aggregates/AggregateNone.h b/src/Aggregates/AggregateNone.h index f7bb0cc..655b282 100644 --- a/src/Aggregates/AggregateNone.h +++ b/src/Aggregates/AggregateNone.h @@ -33,7 +33,7 @@ AggregateNone::AggregateNone(ReactivePredicate condition) template void AggregateNone::OnNext(T value) { - if (!_condition(value)) _state = false; + if (_condition(value)) _state = false; this->_childObservers.OnNext(_state); } diff --git a/src/Aggregates/AggregateRMS.h b/src/Aggregates/AggregateRMS.h index 08c3506..7151558 100644 --- a/src/Aggregates/AggregateRMS.h +++ b/src/Aggregates/AggregateRMS.h @@ -33,7 +33,7 @@ void AggregateRMS::OnNext(T value) { _sumSqr += value * value; _count++; - this->_childObservers.OnNext(sqrt(_sumSqr / _count)); + this->_childObservers.OnNext(static_cast(sqrt(static_cast(_sumSqr) / static_cast(_count)))); } #endif \ No newline at end of file diff --git a/src/Filters/FilterWindowMicros.h b/src/Filters/FilterWindowMicros.h index 363bdc1..179505e 100644 --- a/src/Filters/FilterWindowMicros.h +++ b/src/Filters/FilterWindowMicros.h @@ -32,9 +32,7 @@ void FilterWindowMicros::OnNext(T value) } if (_started && static_cast(micros() - _lastTrigger) <= _interval) - { this->_childObservers.OnNext(value); - } } #endif \ No newline at end of file diff --git a/src/Filters/FilterWindowMillis.h b/src/Filters/FilterWindowMillis.h index 020d061..ee4fbea 100644 --- a/src/Filters/FilterWindowMillis.h +++ b/src/Filters/FilterWindowMillis.h @@ -39,9 +39,7 @@ void FilterWindowMillis::OnNext(T value) } if (_started && static_cast(millis() - _lastTrigger) <= _interval) - { this->_childObservers.OnNext(value); - } } #endif \ No newline at end of file diff --git a/src/Observables/ObservableIntervalMicros.h b/src/Observables/ObservableIntervalMicros.h index 381f2fc..6ea0d81 100644 --- a/src/Observables/ObservableIntervalMicros.h +++ b/src/Observables/ObservableIntervalMicros.h @@ -33,7 +33,6 @@ class ObservableIntervalMicros : public Observable private: bool _isActive; - bool _isExpired; unsigned long _startTime; unsigned long _delay; unsigned long _offset; @@ -121,7 +120,9 @@ unsigned long ObservableIntervalMicros::GetElapsedTime() template unsigned long ObservableIntervalMicros::GetRemainingTime() { - return _interval - micros() + _startTime; + unsigned long elapsed = micros() - _startTime; + if (elapsed >= _interval) return 0; + return _interval - elapsed; } template diff --git a/src/Observables/ObservableIntervalMillis.h b/src/Observables/ObservableIntervalMillis.h index 0315774..5537e55 100644 --- a/src/Observables/ObservableIntervalMillis.h +++ b/src/Observables/ObservableIntervalMillis.h @@ -33,7 +33,6 @@ class ObservableIntervalMillis : public Observable private: bool _isActive; - bool _isExpired; unsigned long _startTime; unsigned long _delay; unsigned long _offset; @@ -121,7 +120,9 @@ unsigned long ObservableIntervalMillis::GetElapsedTime() template unsigned long ObservableIntervalMillis::GetRemainingTime() { - return _interval - millis() + _startTime; + unsigned long elapsed = millis() - _startTime; + if (elapsed >= _interval) return 0; + return _interval - elapsed; } template diff --git a/src/Observables/ObservableRange.h b/src/Observables/ObservableRange.h index 36cf0ee..2d9b3e0 100644 --- a/src/Observables/ObservableRange.h +++ b/src/Observables/ObservableRange.h @@ -54,8 +54,16 @@ void ObservableRange::UnSubscribe(IObserver &observer) template void ObservableRange::Run() { - for (auto i = _start; i <= _end; i += _step) - this->_childObservers.OnNext(i); + if (_step > 0) + { + for (auto i = _start; i <= _end; i += _step) + this->_childObservers.OnNext(i); + } + else if (_step < 0) + { + for (auto i = _start; i >= _end; i += _step) + this->_childObservers.OnNext(i); + } this->_childObservers.OnComplete(); } diff --git a/src/Observables/ObservableRangeDefer.h b/src/Observables/ObservableRangeDefer.h index 09cc17f..4d613f3 100644 --- a/src/Observables/ObservableRangeDefer.h +++ b/src/Observables/ObservableRangeDefer.h @@ -53,13 +53,14 @@ void ObservableRangeDefer::UnSubscribe(IObserver &observer) template void ObservableRangeDefer::Next() { - if (_value > _end) return; + if (_step > 0 && _value > _end) return; + if (_step < 0 && _value < _end) return; T value = _value; this->_childObservers.OnNext(value); _value += _step; - if (_value > _end) + if ((_step > 0 && _value > _end) || (_step < 0 && _value < _end)) this->_childObservers.OnComplete(); } diff --git a/src/Observables/ObservableSerialByte.h b/src/Observables/ObservableSerialByte.h index 0eaefef..1906d4e 100644 --- a/src/Observables/ObservableSerialByte.h +++ b/src/Observables/ObservableSerialByte.h @@ -15,8 +15,8 @@ class ObservableSerial : public Observable { public: ObservableSerial(); - void Subscribe(IObserver &observer); - void UnSubscribe(IObserver &observer); + void Subscribe(IObserver &observer) override; + void UnSubscribe(IObserver &observer) override; void Receive(); private: diff --git a/src/Observables/ObservableSerialDouble.h b/src/Observables/ObservableSerialDouble.h index 8664a10..0820168 100644 --- a/src/Observables/ObservableSerialDouble.h +++ b/src/Observables/ObservableSerialDouble.h @@ -23,7 +23,7 @@ class ObservableSerial : public Observable private: char _separator; - float _data = 0; + double _data = 0; int _dataReal = 0; int _dataDecimal = 0; int _dataPow = 1; diff --git a/src/Observables/ObservableSerialString.h b/src/Observables/ObservableSerialString.h index 6c1ce38..67c501b 100644 --- a/src/Observables/ObservableSerialString.h +++ b/src/Observables/ObservableSerialString.h @@ -47,9 +47,7 @@ inline void ObservableSerial::Receive() { const char newChar = Serial.read(); if (newChar != _separator) - { _buffer.concat(newChar); - } else { _childObservers.OnNext(_buffer); diff --git a/src/Observables/ObservableTimerMicros.h b/src/Observables/ObservableTimerMicros.h index 703b076..e2d35d2 100644 --- a/src/Observables/ObservableTimerMicros.h +++ b/src/Observables/ObservableTimerMicros.h @@ -46,6 +46,7 @@ template ObservableTimerMicros::ObservableTimerMicros(unsigned long interval, unsigned long delay) { _isActive = true; + _isExpired = false; _delay = delay; _offset = delay; _interval = interval; @@ -68,11 +69,13 @@ template void ObservableTimerMicros::Update() { if (_isActive == false) return; + if (_isExpired) return; auto elapsed = static_cast(micros() - _startTime); if (elapsed >= _interval + _offset) { this->_childObservers.OnNext(elapsed); + _isExpired = true; _offset = 0; } } @@ -81,6 +84,7 @@ template void ObservableTimerMicros::Reset() { _isActive = true; + _isExpired = false; _offset = _delay; _startTime = micros(); } @@ -119,7 +123,9 @@ unsigned long ObservableTimerMicros::GetElapsedTime() const template unsigned long ObservableTimerMicros::GetRemainingTime() const { - return _interval - micros() + _startTime; + unsigned long elapsed = micros() - _startTime; + if (elapsed >= _interval) return 0; + return _interval - elapsed; } template diff --git a/src/Observables/ObservableTimerMillis.h b/src/Observables/ObservableTimerMillis.h index 82bc412..9aaa3b5 100644 --- a/src/Observables/ObservableTimerMillis.h +++ b/src/Observables/ObservableTimerMillis.h @@ -46,6 +46,7 @@ template ObservableTimerMillis::ObservableTimerMillis(unsigned long interval, unsigned long delay) { _isActive = true; + _isExpired = false; _delay = delay; _offset = delay; _interval = interval; @@ -68,11 +69,13 @@ template void ObservableTimerMillis::Update() { if (_isActive == false) return; + if (_isExpired) return; auto elapsed = static_cast(millis() - _startTime); if (elapsed >= _interval + _offset) { this->_childObservers.OnNext(elapsed); + _isExpired = true; _offset = 0; } } @@ -81,6 +84,7 @@ template void ObservableTimerMillis::Reset() { _isActive = true; + _isExpired = false; _offset = _delay; _startTime = millis(); } @@ -119,7 +123,9 @@ unsigned long ObservableTimerMillis::GetElapsedTime() const template unsigned long ObservableTimerMillis::GetRemainingTime() const { - return _interval - millis() + _startTime; + unsigned long elapsed = millis() - _startTime; + if (elapsed >= _interval) return 0; + return _interval - elapsed; } template diff --git a/src/Observers/ObserverAnalogOutput.h b/src/Observers/ObserverAnalogOutput.h index d1fc677..2c96f8c 100644 --- a/src/Observers/ObserverAnalogOutput.h +++ b/src/Observers/ObserverAnalogOutput.h @@ -34,7 +34,7 @@ ObserverAnalogOutput::ObserverAnalogOutput(uint8_t pin) template void ObserverAnalogOutput::OnNext(T value) { - analogWrite(value); + analogWrite(_pin, value); } template diff --git a/src/Observers/ObserverDigitalOutput.h b/src/Observers/ObserverDigitalOutput.h index 86ca0c5..adbbbaa 100644 --- a/src/Observers/ObserverDigitalOutput.h +++ b/src/Observers/ObserverDigitalOutput.h @@ -4,12 +4,12 @@ #define _REACTIVEOBSERVERDIGITALOUTPUT_h template -class ObserverDigitalOutput : public IObserver +class ObserverDigitalOutput : public IObserver { public: ObserverDigitalOutput(uint8_t pin); - void OnNext(int value) override; + void OnNext(T value) override; void OnComplete() override; private: @@ -24,7 +24,7 @@ ObserverDigitalOutput::ObserverDigitalOutput(uint8_t pin) } template -void ObserverDigitalOutput::OnNext(int value) +void ObserverDigitalOutput::OnNext(T value) { digitalWrite(_pin, value); } diff --git a/src/Observers/ObserverDoNothing.h b/src/Observers/ObserverDoNothing.h index ac3e76a..c4b6502 100644 --- a/src/Observers/ObserverDoNothing.h +++ b/src/Observers/ObserverDoNothing.h @@ -18,9 +18,6 @@ class ObserverDoNothing : public IObserver void OnNext(T value) override; void OnComplete() override; - -private: - ReactiveAction _doAction; }; template diff --git a/src/Observers/ObserverFinally.h b/src/Observers/ObserverFinally.h index 8d36e74..d48d723 100644 --- a/src/Observers/ObserverFinally.h +++ b/src/Observers/ObserverFinally.h @@ -38,7 +38,7 @@ void ObserverFinally::OnNext(T value) template void ObserverFinally::OnComplete() { - _action(); + if (_action != nullptr) _action(); } #endif \ No newline at end of file diff --git a/src/Observers/ObserverSerial.h b/src/Observers/ObserverSerial.h index 88813bf..2fcb8f4 100644 --- a/src/Observers/ObserverSerial.h +++ b/src/Observers/ObserverSerial.h @@ -14,9 +14,6 @@ template class ObserverSerial : public IObserver { public: - - -private: void OnNext(T value) override; void OnComplete() override; }; diff --git a/src/Operators/OperatorBatch.h b/src/Operators/OperatorBatch.h index 53d07f8..98c9db6 100644 --- a/src/Operators/OperatorBatch.h +++ b/src/Operators/OperatorBatch.h @@ -32,16 +32,14 @@ OperatorBatch::OperatorBatch(size_t N) template void OperatorBatch::OnNext(T value) { - if (_index < _num_elements) - { - this->_childObservers.OnNext(value); - _index++; - } - else + if (_index >= _num_elements) { _index = 0; this->_childObservers.OnComplete(); } + + this->_childObservers.OnNext(value); + _index++; } #endif \ No newline at end of file diff --git a/src/Operators/OperatorFirst.h b/src/Operators/OperatorFirst.h index 838d175..b822ccc 100644 --- a/src/Operators/OperatorFirst.h +++ b/src/Operators/OperatorFirst.h @@ -45,6 +45,8 @@ void OperatorFirst::OnComplete() { if (_any) this->_childObservers.OnNext(_first); + + this->_childObservers.OnComplete(); } #endif \ No newline at end of file diff --git a/src/Operators/OperatorForEach.h b/src/Operators/OperatorForEach.h index fbacfd4..104c9ba 100644 --- a/src/Operators/OperatorForEach.h +++ b/src/Operators/OperatorForEach.h @@ -40,6 +40,7 @@ void OperatorForEach::OnNext(T value) template inline void OperatorForEach::OnComplete() { + this->_childObservers.OnComplete(); } #endif \ No newline at end of file diff --git a/src/Operators/OperatorRepeat.h b/src/Operators/OperatorRepeat.h index 8420d18..3ffc4fa 100644 --- a/src/Operators/OperatorRepeat.h +++ b/src/Operators/OperatorRepeat.h @@ -38,15 +38,14 @@ void OperatorRepeat::OnNext(T value) template void OperatorRepeat::OnComplete() { - _repetition--; + if (_repetition > 0) _repetition--; + if (_repetition > 0) { if (this->_parentObservable != nullptr) this->_parentObservable->Reset(); } else - { this->_childObservers.OnComplete(); - } } #endif \ No newline at end of file diff --git a/src/Operators/OperatorSkip.h b/src/Operators/OperatorSkip.h index 2a84db2..29ba640 100644 --- a/src/Operators/OperatorSkip.h +++ b/src/Operators/OperatorSkip.h @@ -21,7 +21,6 @@ class OperatorSkip : public Operator private: size_t _index = 0; size_t _num_elements = 0; - bool _completed = false; }; template diff --git a/src/Operators/OperatorTake.h b/src/Operators/OperatorTake.h index b2a5af5..c64e613 100644 --- a/src/Operators/OperatorTake.h +++ b/src/Operators/OperatorTake.h @@ -35,17 +35,21 @@ void OperatorTake::OnNext(T value) { if (_completed) return; - if (_index < _num_elements) - { - this->_childObservers.OnNext(value); - } - else + if (_num_elements == 0) { this->_childObservers.OnComplete(); _completed = true; + return; } + this->_childObservers.OnNext(value); _index++; + + if (_index >= _num_elements) + { + this->_childObservers.OnComplete(); + _completed = true; + } } #endif \ No newline at end of file diff --git a/src/Operators/OperatorTakeFirst.h b/src/Operators/OperatorTakeFirst.h index 017082a..c977632 100644 --- a/src/Operators/OperatorTakeFirst.h +++ b/src/Operators/OperatorTakeFirst.h @@ -33,6 +33,7 @@ void OperatorTakeFirst::OnNext(T value) if (_completed) return; this->_childObservers.OnNext(value); + this->_childObservers.OnComplete(); _completed = true; } diff --git a/src/Operators/OperatorTakeLast.h b/src/Operators/OperatorTakeLast.h index 36ba448..6906ffa 100644 --- a/src/Operators/OperatorTakeLast.h +++ b/src/Operators/OperatorTakeLast.h @@ -17,6 +17,11 @@ class OperatorTakeLast : public Operator OperatorTakeLast(); void OnNext(T value) override; + void OnComplete() override; + +private: + T _last = T(); + bool _any = false; }; template @@ -27,7 +32,17 @@ OperatorTakeLast::OperatorTakeLast() template void OperatorTakeLast::OnNext(T value) { - this->_childObservers.OnNext(value); + _last = value; + _any = true; +} + +template +void OperatorTakeLast::OnComplete() +{ + if (!_any) return; + + this->_childObservers.OnNext(_last); + this->_childObservers.OnComplete(); } #endif \ No newline at end of file diff --git a/src/Operators/OperatorTakeUntil.h b/src/Operators/OperatorTakeUntil.h index dfbd10b..60d057c 100644 --- a/src/Operators/OperatorTakeUntil.h +++ b/src/Operators/OperatorTakeUntil.h @@ -36,9 +36,7 @@ void OperatorTakeUntil::OnNext(T value) if (_completed) return; if (!this->_condition(value)) - { this->_childObservers.OnNext(value); - } else { _completed = true; diff --git a/src/Operators/OperatorTakeWhile.h b/src/Operators/OperatorTakeWhile.h index e5782a7..21f90e3 100644 --- a/src/Operators/OperatorTakeWhile.h +++ b/src/Operators/OperatorTakeWhile.h @@ -36,9 +36,7 @@ void OperatorTakeWhile::OnNext(T value) if (_completed) return; if (this->_condition(value)) - { this->_childObservers.OnNext(value); - } else { _completed = true; diff --git a/src/Operators/OperatorTimeoutMicros.h b/src/Operators/OperatorTimeoutMicros.h index 6f5ea9d..36e006f 100644 --- a/src/Operators/OperatorTimeoutMicros.h +++ b/src/Operators/OperatorTimeoutMicros.h @@ -14,32 +14,31 @@ template class OperatorTimeoutMicros : public Operator { public: - OperatorTimeoutMicros(unsigned long interval, ReactiveAction action); + OperatorTimeoutMicros(unsigned long interval, ReactiveCallback action); void OnNext(T value) override; void OnComplete() override; void Update(); private: - ReactiveAction _doAction; - unsigned long _starTime; + ReactiveCallback _doAction; + unsigned long _startTime; unsigned long _interval; bool _completed = false; }; template -OperatorTimeoutMicros::OperatorTimeoutMicros(unsigned long interval, ReactiveAction action) +OperatorTimeoutMicros::OperatorTimeoutMicros(unsigned long interval, ReactiveCallback action) { _doAction = action; _interval = interval; - _starTime = micros(); + _startTime = micros(); } template void OperatorTimeoutMicros::OnNext(T value) { - _doAction(value); - _starTime = micros(); + _startTime = micros(); this->_childObservers.OnNext(value); } @@ -56,7 +55,7 @@ inline void OperatorTimeoutMicros::Update() { if (_completed) return; - if (millis() - _starTime > _interval) + if (micros() - _startTime > _interval) { if (_doAction != nullptr) _doAction(); _completed = true; diff --git a/src/Operators/OperatorTimeoutMillis.h b/src/Operators/OperatorTimeoutMillis.h index d90e87e..58709fc 100644 --- a/src/Operators/OperatorTimeoutMillis.h +++ b/src/Operators/OperatorTimeoutMillis.h @@ -14,32 +14,31 @@ template class OperatorTimeoutMillis : public Operator { public: - OperatorTimeoutMillis(unsigned long interval, ReactiveAction action); + OperatorTimeoutMillis(unsigned long interval, ReactiveCallback action); void OnNext(T value) override; void OnComplete() override; void Update(); private: - ReactiveAction _doAction; - unsigned long _starTime; + ReactiveCallback _doAction; + unsigned long _startTime; unsigned long _interval; bool _completed = false; }; template -OperatorTimeoutMillis::OperatorTimeoutMillis(unsigned long interval, ReactiveAction action) +OperatorTimeoutMillis::OperatorTimeoutMillis(unsigned long interval, ReactiveCallback action) { _doAction = action; _interval = interval; - _starTime = millis(); + _startTime = millis(); } template void OperatorTimeoutMillis::OnNext(T value) { - _doAction(value); - _starTime = millis(); + _startTime = millis(); this->_childObservers.OnNext(value); } @@ -56,7 +55,7 @@ inline void OperatorTimeoutMillis::Update() { if (_completed) return; - if (millis() - _starTime > _interval) + if (millis() - _startTime > _interval) { if (_doAction != nullptr) _doAction(); _completed = true; diff --git a/src/ReactiveArduinoCore.h b/src/ReactiveArduinoCore.h index 37e3ae6..438c992 100644 --- a/src/ReactiveArduinoCore.h +++ b/src/ReactiveArduinoCore.h @@ -168,8 +168,8 @@ class Observable : IObservable, IResetable OperatorBatch& Batch(size_t num); OperatorIf& If(ReactivePredicate condition, ReactiveAction action); OperatorForEach& ForEach(ReactiveAction action); - OperatorTimeoutMillis& TimeoutMillis(ReactiveAction action); - OperatorTimeoutMicros& TimeoutMicros(ReactiveAction action); + OperatorTimeoutMillis& TimeoutMillis(unsigned long interval, ReactiveCallback action); + OperatorTimeoutMicros& TimeoutMicros(unsigned long interval, ReactiveCallback action); OperatorReset& DoReset(); OperatorNoReset& NotReset(); OperatorLoop& Loop(); @@ -190,10 +190,10 @@ class Observable : IObservable, IResetable TransformationElapsedMillis& ElapsedMillis(); TransformationElapsedMicros& ElapsedMicros(); TransformationFrequency& Frequency(); - TransformationThreshold& Threshold(T threshold, int state = LOW); - TransformationThreshold& DoubleThreshold(T lowThreshold, T highThreshold, int state = LOW); - TransformationToggle& Toggle(int state = LOW); - TransformationAdcToVoltage& AdcToVoltage(T input_max = 1023, T output_max = 5.0); + TransformationThreshold& Threshold(T threshold, bool state = false); + TransformationThreshold& DoubleThreshold(T lowThreshold, T highThreshold, bool state = false); + TransformationToggle& Toggle(bool state = false); + TransformationAdcToVoltage& AdcToVoltage(float input_max = 1023.0f, float output_max = 5.0f); TransformationSplit& Split(char separator = ','); TransformationJoin& Join(char separator = ','); TransformationStringBuffer & StringBuffer(); @@ -237,7 +237,7 @@ class Observable : IObservable, IResetable AggregateAll& All(ReactivePredicate condition); AggregateNone& None(ReactivePredicate condition); - ObserverSerial ToSerial(); + ObserverSerial& ToSerial(); ObserverDo& Do(ReactiveAction action); ObserverFinally& Finally(ReactiveCallback action); ObserverDoAndFinally& DoAndFinally(ReactiveAction doAction, ReactiveCallback finallyAction); @@ -385,17 +385,17 @@ auto Observable::ForEach(ReactiveAction action) -> OperatorForEach& } template -auto Observable::TimeoutMillis(ReactiveAction action) -> OperatorTimeoutMillis& +auto Observable::TimeoutMillis(unsigned long interval, ReactiveCallback action) -> OperatorTimeoutMillis& { - auto newOp = new OperatorTimeoutMillis(action); + auto newOp = new OperatorTimeoutMillis(interval, action); Compound(*this, *newOp); return *newOp; } template -auto Observable::TimeoutMicros(ReactiveAction action) -> OperatorTimeoutMicros& +auto Observable::TimeoutMicros(unsigned long interval, ReactiveCallback action) -> OperatorTimeoutMicros& { - auto newOp = new OperatorTimeoutMicros(action); + auto newOp = new OperatorTimeoutMicros(interval, action); Compound(*this, *newOp); return *newOp; } @@ -560,7 +560,7 @@ auto Observable::Frequency() -> TransformationFrequency& } template -auto Observable::Threshold(T threshold, int state) -> TransformationThreshold& +auto Observable::Threshold(T threshold, bool state) -> TransformationThreshold& { auto newOp = new TransformationThreshold(threshold, state); Compound(*this, *newOp); @@ -568,7 +568,7 @@ auto Observable::Threshold(T threshold, int state) -> TransformationThreshold } template -auto Observable::DoubleThreshold(T lowThreshold, T highThreshold, int state) -> TransformationThreshold& +auto Observable::DoubleThreshold(T lowThreshold, T highThreshold, bool state) -> TransformationThreshold& { auto newOp = new TransformationThreshold(lowThreshold, highThreshold, state); Compound(*this, *newOp); @@ -576,7 +576,7 @@ auto Observable::DoubleThreshold(T lowThreshold, T highThreshold, int state) } template -auto Observable::Toggle(int state) -> TransformationToggle& +auto Observable::Toggle(bool state) -> TransformationToggle& { auto newOp = new TransformationToggle(state); Compound(*this, *newOp); @@ -584,7 +584,7 @@ auto Observable::Toggle(int state) -> TransformationToggle& } template -auto Observable::AdcToVoltage(T input_max, T output_max) -> TransformationAdcToVoltage& +auto Observable::AdcToVoltage(float input_max, float output_max) -> TransformationAdcToVoltage& { auto newOp = new TransformationAdcToVoltage(input_max, output_max); Compound(*this, *newOp); @@ -912,7 +912,7 @@ auto Observable::None(ReactivePredicate condition) -> AggregateNone& } template -auto Observable::ToSerial() -> ObserverSerial +auto Observable::ToSerial() -> ObserverSerial& { auto newOp = new ObserverSerial(); Subscribe(*newOp); diff --git a/src/ReactiveArduinoLib.h b/src/ReactiveArduinoLib.h index 272cbcb..d88da47 100644 --- a/src/ReactiveArduinoLib.h +++ b/src/ReactiveArduinoLib.h @@ -111,7 +111,7 @@ namespace Reactive template auto FromSerial(char separator) -> ObservableSerial& { - return *(new ObservableSerial()); + return *(new ObservableSerial(separator)); } template <> @@ -314,15 +314,15 @@ namespace Reactive } template - OperatorTimeoutMillis& TimeoutMillis(ReactiveAction action) + OperatorTimeoutMillis& TimeoutMillis(unsigned long interval, ReactiveCallback action) { - return *(new OperatorTimeoutMillis(action)); + return *(new OperatorTimeoutMillis(interval, action)); } template - OperatorTimeoutMicros& TimeoutMicros(ReactiveAction action) + OperatorTimeoutMicros& TimeoutMicros(unsigned long interval, ReactiveCallback action) { - return *(new OperatorTimeoutMicros(action)); + return *(new OperatorTimeoutMicros(interval, action)); } template @@ -457,7 +457,7 @@ namespace Reactive template TransformationThreshold& DoubleThreshold(T lowThreshold, T highThreshold) { - return *(new TransformationThreshold(lowThreshold, highThreshold)); + return *(new TransformationThreshold(lowThreshold, highThreshold, false)); } template @@ -473,7 +473,7 @@ namespace Reactive } template - TransformationAdcToVoltage& AdcToVoltage(T input_max = 1023, T output_max = 5.0) + TransformationAdcToVoltage& AdcToVoltage(float input_max = 1023.0f, float output_max = 5.0f) { return *(new TransformationAdcToVoltage(input_max, output_max)); } diff --git a/src/Transformations/TransformationAdcToVoltage.h b/src/Transformations/TransformationAdcToVoltage.h index fa9d0f3..785ee99 100644 --- a/src/Transformations/TransformationAdcToVoltage.h +++ b/src/Transformations/TransformationAdcToVoltage.h @@ -14,17 +14,17 @@ template class TransformationAdcToVoltage : public Operator { public: - TransformationAdcToVoltage(T input_max = 1023, T output_max = 5.0); + TransformationAdcToVoltage(float input_max = 1023.0f, float output_max = 5.0f); void OnNext(T value) override; private: - T _input_max = T(); - T _output_max = T(); + float _input_max = 0.0f; + float _output_max = 0.0f; }; template -TransformationAdcToVoltage::TransformationAdcToVoltage(T input_max, T output_max) +TransformationAdcToVoltage::TransformationAdcToVoltage(float input_max, float output_max) { _input_max = input_max; _output_max = output_max; @@ -34,7 +34,7 @@ TransformationAdcToVoltage::TransformationAdcToVoltage(T input_max, T output_ template void TransformationAdcToVoltage::OnNext(T value) { - this->_childObservers.OnNext((value * _output_max) / _input_max); + this->_childObservers.OnNext((static_cast(value) * _output_max) / _input_max); } #endif \ No newline at end of file diff --git a/src/Transformations/TransformationElapsedMicros.h b/src/Transformations/TransformationElapsedMicros.h index 4a83a99..13e86ee 100644 --- a/src/Transformations/TransformationElapsedMicros.h +++ b/src/Transformations/TransformationElapsedMicros.h @@ -21,26 +21,26 @@ class TransformationElapsedMicros : public Operator void Reset() override; private: - unsigned long _starTime; + unsigned long _startTime; }; template TransformationElapsedMicros::TransformationElapsedMicros() { - _starTime = micros(); + _startTime = micros(); } template void TransformationElapsedMicros::Reset() { - _starTime = micros(); + _startTime = micros(); } template void TransformationElapsedMicros::OnNext(T value) { - this->_childObservers.OnNext(micros() - _starTime); - _starTime = micros(); + this->_childObservers.OnNext(micros() - _startTime); + _startTime = micros(); } #endif diff --git a/src/Transformations/TransformationElapsedMillis.h b/src/Transformations/TransformationElapsedMillis.h index 760b4de..494e301 100644 --- a/src/Transformations/TransformationElapsedMillis.h +++ b/src/Transformations/TransformationElapsedMillis.h @@ -21,26 +21,26 @@ class TransformationElapsedMillis : public Operator void Reset() override; private: - unsigned long _starTime; + unsigned long _startTime; }; template TransformationElapsedMillis::TransformationElapsedMillis() { - _starTime = millis(); + _startTime = millis(); } template void TransformationElapsedMillis::Reset() { - _starTime = millis(); + _startTime = millis(); } template void TransformationElapsedMillis::OnNext(T value) { - this->_childObservers.OnNext(millis() - _starTime); - _starTime = millis(); + this->_childObservers.OnNext(millis() - _startTime); + _startTime = millis(); } #endif \ No newline at end of file diff --git a/src/Transformations/TransformationFrequency.h b/src/Transformations/TransformationFrequency.h index 0229fc6..b2c415b 100644 --- a/src/Transformations/TransformationFrequency.h +++ b/src/Transformations/TransformationFrequency.h @@ -27,25 +27,26 @@ class TransformationFrequency : public Operator void Reset() override; private: - unsigned long _starTime; + unsigned long _startTime; }; template TransformationFrequency::TransformationFrequency() { - _starTime = millis(); + _startTime = millis(); } template void TransformationFrequency::Reset() { - _starTime = millis(); + _startTime = millis(); } template void TransformationFrequency::OnNext(T value) { - this->_childObservers.OnNext(1000.0 / (millis() - _starTime)); - _starTime = millis(); + unsigned long elapsed = millis() - _startTime; + this->_childObservers.OnNext(elapsed > 0 ? 1000.0f / elapsed : 0.0f); + _startTime = millis(); } #endif diff --git a/src/Transformations/TransformationJoin.h b/src/Transformations/TransformationJoin.h index 6ef4067..2569ccd 100644 --- a/src/Transformations/TransformationJoin.h +++ b/src/Transformations/TransformationJoin.h @@ -40,9 +40,7 @@ void TransformationJoin::OnNext(T value) _isFirst = false; } else - { _buffer = _buffer + _separator + String(value); - } this->_childObservers.OnNext(_buffer); } diff --git a/src/Transformations/TransformationScale.h b/src/Transformations/TransformationScale.h index 7ffa872..c08e11e 100644 --- a/src/Transformations/TransformationScale.h +++ b/src/Transformations/TransformationScale.h @@ -38,6 +38,12 @@ TransformationScale::TransformationScale(T input_min, T input_max, T output_m template void TransformationScale::OnNext(T value) { + if (_input_max == _input_min) + { + this->_childObservers.OnNext(_output_min); + return; + } + T scaled = (value - _input_min) * (_output_max - _output_min) / (_input_max - _input_min) + _output_min; this->_childObservers.OnNext(scaled); } diff --git a/src/Transformations/TransformationSplit.h b/src/Transformations/TransformationSplit.h index 804d584..3d19cf9 100644 --- a/src/Transformations/TransformationSplit.h +++ b/src/Transformations/TransformationSplit.h @@ -14,8 +14,6 @@ template class TransformationSplit : public Operator { public: - ReactiveFunction _function; - TransformationSplit(char separator = ','); void OnNext(T value) override; @@ -34,15 +32,13 @@ TransformationSplit::TransformationSplit(char separator) template void TransformationSplit::OnNext(T value) { - size_t counter = 0; size_t lastIndex = 0; _buffer = value; for (size_t index = 0; index < _buffer.length(); index++) { - if (_buffer.substring(index, index + 1) == ",") + if (_buffer.charAt(index) == _separator) { this->_childObservers.OnNext(_buffer.substring(lastIndex, index)); lastIndex = index + 1; - counter++; } if (index == _buffer.length() - 1) diff --git a/src/Transformations/TransformationThreshold.h b/src/Transformations/TransformationThreshold.h index b83ebc5..65a023d 100644 --- a/src/Transformations/TransformationThreshold.h +++ b/src/Transformations/TransformationThreshold.h @@ -14,21 +14,20 @@ template class TransformationThreshold : public Operator { public: - TransformationThreshold(T threshold) : TransformationThreshold(threshold, threshold, LOW) {} - TransformationThreshold(T threshold, int state) : TransformationThreshold(threshold, threshold, state) {} - TransformationThreshold(T lowThreshold, T highThreshold) : TransformationThreshold(lowThreshold, highThreshold, LOW) {} - TransformationThreshold(T lowThreshold, T highThreshold, int state); + TransformationThreshold(T threshold) : TransformationThreshold(threshold, threshold, false) {} + TransformationThreshold(T threshold, bool state) : TransformationThreshold(threshold, threshold, state) {} + TransformationThreshold(T lowThreshold, T highThreshold, bool state); void OnNext(T value) override; private: T _fallThreshold = T(); T _riseThreshold = T(); - int _state; + bool _state = false; }; template -TransformationThreshold::TransformationThreshold(T threshold1, T threshold2, int state) +TransformationThreshold::TransformationThreshold(T threshold1, T threshold2, bool state) { _fallThreshold = threshold1 <= threshold2 ? threshold1 : threshold2; _riseThreshold = threshold1 > threshold2 ? threshold1 : threshold2; @@ -38,15 +37,11 @@ TransformationThreshold::TransformationThreshold(T threshold1, T threshold2, template void TransformationThreshold::OnNext(T value) { - if (_state == LOW && value > _riseThreshold) - { - _state = HIGH; - } - - if (_state == HIGH && value < _fallThreshold) - { - _state = LOW; - } + if (!_state && value > _riseThreshold) + _state = true; + + if (_state && value < _fallThreshold) + _state = false; this->_childObservers.OnNext(_state); } diff --git a/src/Transformations/TransformationTimestampMicros.h b/src/Transformations/TransformationTimestampMicros.h index b9333a1..fd69ff1 100644 --- a/src/Transformations/TransformationTimestampMicros.h +++ b/src/Transformations/TransformationTimestampMicros.h @@ -21,24 +21,24 @@ class TransformationTimestampMicros : public Operator void Reset() override; private: - unsigned long _starTime; + unsigned long _startTime; }; template TransformationTimestampMicros::TransformationTimestampMicros() { - _starTime = micros(); + _startTime = micros(); } template void TransformationTimestampMicros::Reset() { - _starTime = micros(); + _startTime = micros(); } template void TransformationTimestampMicros::OnNext(T value) { - this->_childObservers.OnNext(micros() - _starTime); + this->_childObservers.OnNext(micros() - _startTime); } #endif \ No newline at end of file diff --git a/src/Transformations/TransformationTimestampMillis.h b/src/Transformations/TransformationTimestampMillis.h index 26e0775..43b375f 100644 --- a/src/Transformations/TransformationTimestampMillis.h +++ b/src/Transformations/TransformationTimestampMillis.h @@ -21,24 +21,24 @@ class TransformationTimestampMillis : public Operator void Reset() override; private: - unsigned long _starTime; + unsigned long _startTime; }; template TransformationTimestampMillis::TransformationTimestampMillis() { - _starTime = millis(); + _startTime = millis(); } template void TransformationTimestampMillis::Reset() { - _starTime = millis(); + _startTime = millis(); } template void TransformationTimestampMillis::OnNext(T value) { - this->_childObservers.OnNext(millis() - _starTime); + this->_childObservers.OnNext(millis() - _startTime); } #endif \ No newline at end of file diff --git a/src/Transformations/TransformationToggle.h b/src/Transformations/TransformationToggle.h index 4c9d969..c73baf4 100644 --- a/src/Transformations/TransformationToggle.h +++ b/src/Transformations/TransformationToggle.h @@ -14,18 +14,18 @@ template class TransformationToggle : public Operator { public: - TransformationToggle(int state = LOW); + TransformationToggle(bool state = false); void OnNext(T value) override; private: - int _state = false; + bool _state = false; }; template -TransformationToggle::TransformationToggle(int state) +TransformationToggle::TransformationToggle(bool state) { - _state = false; + _state = state; } template diff --git a/test/TestHelpers.h b/test/TestHelpers.h new file mode 100644 index 0000000..fd24501 --- /dev/null +++ b/test/TestHelpers.h @@ -0,0 +1,35 @@ +// Shared helpers for the ReactiveArduino unit tests. +#ifndef TEST_HELPERS_H +#define TEST_HELPERS_H + +#include "WProgram.h" +#include + +// A terminal observer that records every value and completion it receives. +template +struct Sink : IObserver +{ + std::vector values; + int completeCount = 0; + + void OnNext(T v) override { values.push_back(v); } + void OnComplete() override { completeCount++; } +}; + +// Reset the Arduino mock state before each test. +inline void resetMocks() +{ + g_millis = 0; + g_micros = 0; + g_pinModePin = 0xFF; + g_pinModeMode = 0xFF; + g_digitalPin = 0xFF; + g_digitalValue = 0xFF; + g_analogPin = 0xFF; + g_analogValue = -1; + g_analogRead = 0; + g_digitalRead = 0; + Serial.printCount = 0; +} + +#endif diff --git a/test/WProgram.h b/test/WProgram.h new file mode 100644 index 0000000..af13a97 --- /dev/null +++ b/test/WProgram.h @@ -0,0 +1,107 @@ +// Minimal Arduino API stub used by the native PlatformIO unit tests. +// ReactiveArduinoLib.h includes "WProgram.h" when ARDUINO is not defined, +// so this file is picked up from the `test` include path. +#ifndef WPROGRAM_STUB_H +#define WPROGRAM_STUB_H + +#include +#include +#include +#include +#include +#include +#include + +typedef unsigned char byte; + +// ---- Minimal Arduino String ---- +class String +{ +public: + String() {} + String(const char* s) { if (s) _s = s; } + String(char c) { _s.assign(1, c); } + String(int v) { char b[16]; snprintf(b, sizeof(b), "%d", v); _s = b; } + String(unsigned int v) { char b[16]; snprintf(b, sizeof(b), "%u", v); _s = b; } + String(long v) { char b[24]; snprintf(b, sizeof(b), "%ld", v); _s = b; } + String(unsigned long v) { char b[24]; snprintf(b, sizeof(b), "%lu", v); _s = b; } + String(float v) { char b[32]; snprintf(b, sizeof(b), "%f", v); _s = b; } + String(double v) { char b[40]; snprintf(b, sizeof(b), "%f", v); _s = b; } + + const char* c_str() const { return _s.c_str(); } + unsigned int length() const { return (unsigned int)_s.size(); } + char charAt(unsigned int i) const { return (i < _s.size()) ? _s[i] : '\0'; } + + String substring(unsigned int from, unsigned int to) const + { + if (from >= _s.size()) return String(); + if (to > _s.size()) to = (unsigned int)_s.size(); + if (to < from) to = from; + return String(_s.substr(from, to - from).c_str()); + } + + void concat(char c) { _s += c; } + void concat(const char* s) { if (s) _s += s; } + int toInt() const { return (int)atoi(_s.c_str()); } + float toFloat() const { return (float)atof(_s.c_str()); } + + String operator+(const String& o) const { String r(_s.c_str()); r._s += o._s; return r; } + String& operator=(const String& o) { _s = o._s; return *this; } + String& operator=(const char* o) { _s = (o ? o : ""); return *this; } + bool operator==(const String& o) const { return _s == o._s; } + bool operator==(const char* o) const { return _s == (o ? o : ""); } + +private: + std::string _s; +}; + +// ---- Controllable time ---- +inline unsigned long g_millis = 0; +inline unsigned long g_micros = 0; +inline unsigned long millis() { return g_millis; } +inline unsigned long micros() { return g_micros; } + +// ---- Pin I/O capture ---- +inline uint8_t g_pinModePin = 0xFF; +inline uint8_t g_pinModeMode = 0xFF; +inline void pinMode(uint8_t pin, uint8_t mode) { g_pinModePin = pin; g_pinModeMode = mode; } + +inline uint8_t g_digitalPin = 0xFF; +inline uint8_t g_digitalValue = 0xFF; +inline void digitalWrite(uint8_t pin, uint8_t value) { g_digitalPin = pin; g_digitalValue = value; } + +inline uint8_t g_analogPin = 0xFF; +inline int g_analogValue = -1; +inline void analogWrite(uint8_t pin, int value) { g_analogPin = pin; g_analogValue = value; } + +inline int g_analogRead = 0; +inline int analogRead(uint8_t) { return g_analogRead; } + +inline int g_digitalRead = 0; +inline int digitalRead(uint8_t) { return g_digitalRead; } + +// ---- Serial capture ---- +struct SerialStub +{ + int printCount = 0; + int available() { return 0; } + int read() { return -1; } + void begin(unsigned long) {} + template void println(const U&) { printCount++; } +}; +inline SerialStub Serial; + +#ifndef LOW +#define LOW 0 +#endif +#ifndef HIGH +#define HIGH 1 +#endif +#ifndef INPUT +#define INPUT 0 +#endif +#ifndef OUTPUT +#define OUTPUT 1 +#endif + +#endif diff --git a/test/test_aggregates.cpp b/test/test_aggregates.cpp new file mode 100644 index 0000000..00bc470 --- /dev/null +++ b/test/test_aggregates.cpp @@ -0,0 +1,128 @@ +#include +#include "ReactiveArduinoLib.h" +using namespace Reactive; +#include "TestHelpers.h" + +static bool isPositive(int v) { return v > 0; } + +void test_count(void) +{ + int arr[4] = {10, 20, 30, 40}; + ObservableArray src(arr, 4); + auto& agg = src.Count(); + Sink s; + agg.Subscribe(s); + TEST_ASSERT_EQUAL_INT(4, s.values.size()); + TEST_ASSERT_EQUAL_INT(4, s.values.back()); +} + +void test_countdown(void) +{ + int arr[4] = {0, 0, 0, 0}; + ObservableArray src(arr, 4); + auto& agg = src.CountDown(3); + Sink s; + agg.Subscribe(s); + TEST_ASSERT_EQUAL_INT(3, s.values.size()); // 2, 1, 0 + TEST_ASSERT_EQUAL_INT(0, s.values.back()); +} + +void test_sum(void) +{ + int arr[4] = {1, 2, 3, 4}; + ObservableArray src(arr, 4); + auto& agg = src.Sum(); + Sink s; + agg.Subscribe(s); + TEST_ASSERT_EQUAL_INT(10, s.values.back()); +} + +void test_min(void) +{ + int arr[5] = {5, 2, 8, 1, 3}; + ObservableArray src(arr, 5); + auto& agg = src.Min(); + Sink s; + agg.Subscribe(s); + TEST_ASSERT_EQUAL_INT(1, s.values.back()); +} + +void test_max(void) +{ + int arr[5] = {5, 2, 8, 1, 3}; + ObservableArray src(arr, 5); + auto& agg = src.Max(); + Sink s; + agg.Subscribe(s); + TEST_ASSERT_EQUAL_INT(8, s.values.back()); +} + +void test_average_float(void) +{ + float arr[4] = {1.0f, 2.0f, 3.0f, 4.0f}; + ObservableArray src(arr, 4); + auto& agg = src.Average(); + Sink s; + agg.Subscribe(s); + TEST_ASSERT_FLOAT_WITHIN(0.0001f, 2.5f, s.values.back()); +} + +void test_rms_float(void) +{ + float arr[2] = {3.0f, 4.0f}; + ObservableArray src(arr, 2); + auto& agg = src.RMS(); + Sink s; + agg.Subscribe(s); + TEST_ASSERT_FLOAT_WITHIN(0.0001f, 3.5355339f, s.values.back()); +} + +void test_any(void) +{ + int arr[3] = {1, 2, 3}; + ObservableArray src(arr, 3); + auto& agg = src.Any(isPositive); + Sink s; + agg.Subscribe(s); + TEST_ASSERT_TRUE(s.values.back()); +} + +void test_all_true(void) +{ + int arr[3] = {1, 2, 3}; + ObservableArray src(arr, 3); + auto& agg = src.All(isPositive); + Sink s; + agg.Subscribe(s); + TEST_ASSERT_TRUE(s.values.back()); +} + +void test_all_detects_failure(void) +{ + int arr[3] = {1, -2, 3}; + ObservableArray src(arr, 3); + auto& agg = src.All(isPositive); + Sink s; + agg.Subscribe(s); + TEST_ASSERT_FALSE(s.values.back()); +} + +void test_none(void) +{ + int arr[3] = {1, 2, 3}; + ObservableArray src(arr, 3); + auto& agg = src.None(isPositive); + Sink s; + agg.Subscribe(s); + TEST_ASSERT_FALSE(s.values.back()); +} + +void test_none_detects_match(void) +{ + int arr[3] = {-1, 2, -3}; + ObservableArray src(arr, 3); + auto& agg = src.None(isPositive); + Sink s; + agg.Subscribe(s); + TEST_ASSERT_FALSE(s.values.back()); +} diff --git a/test/test_filters.cpp b/test/test_filters.cpp new file mode 100644 index 0000000..d5dc485 --- /dev/null +++ b/test/test_filters.cpp @@ -0,0 +1,170 @@ +#include +#include +#include "ReactiveArduinoLib.h" +using namespace Reactive; +#include "TestHelpers.h" + +void test_median3_bruteforce(void) +{ + int p[3] = {0, 1, 2}; + do + { + ObservableArray src(p, 3); + auto& f = src.Median3(); + Sink s; + f.Subscribe(s); + TEST_ASSERT_EQUAL_INT(1, s.values.back()); + } while (std::next_permutation(p, p + 3)); +} + +void test_median5_bruteforce(void) +{ + int p[5] = {0, 1, 2, 3, 4}; + do + { + ObservableArray src(p, 5); + auto& f = src.Median5(); + Sink s; + f.Subscribe(s); + TEST_ASSERT_EQUAL_INT(2, s.values.back()); + } while (std::next_permutation(p, p + 5)); +} + +void test_moving_average(void) +{ + float seq[5] = {1.0f, 2.0f, 3.0f, 4.0f, 5.0f}; + ObservableArray src(seq, 5); + auto& f = src.MovingAverage(3); + Sink s; + f.Subscribe(s); + TEST_ASSERT_EQUAL_INT(5, s.values.size()); + TEST_ASSERT_FLOAT_WITHIN(0.0001f, 2.0f, s.values[2]); // (1+2+3)/3 + TEST_ASSERT_FLOAT_WITHIN(0.0001f, 4.0f, s.values[4]); // (3+4+5)/3 +} + +void test_moving_rms(void) +{ + float seq[2] = {3.0f, 4.0f}; + ObservableArray src(seq, 2); + auto& f = src.MovingRMS(2); + Sink s; + f.Subscribe(s); + TEST_ASSERT_FLOAT_WITHIN(0.0001f, 3.5355339f, s.values.back()); // sqrt((9+16)/2) +} + +void test_on_rising(void) +{ + int arr[4] = {1, 2, 1, 0}; + ObservableArray src(arr, 4); + auto& f = src.OnRising(); + Sink s; + f.Subscribe(s); + TEST_ASSERT_EQUAL_INT(1, s.values.size()); + TEST_ASSERT_EQUAL_INT(2, s.values[0]); +} + +void test_on_falling(void) +{ + int arr[4] = {1, 2, 1, 0}; + ObservableArray src(arr, 4); + auto& f = src.OnFalling(); + Sink s; + f.Subscribe(s); + TEST_ASSERT_EQUAL_INT(2, s.values.size()); + TEST_ASSERT_EQUAL_INT(1, s.values[0]); + TEST_ASSERT_EQUAL_INT(0, s.values[1]); +} + +void test_low_pass(void) +{ + float arr[1] = {1.0f}; + ObservableArray src(arr, 1); + auto& f = src.LowPass(0.5); + Sink s; + f.Subscribe(s); + TEST_ASSERT_FLOAT_WITHIN(0.0001f, 0.5f, s.values[0]); +} + +void test_high_pass(void) +{ + float arr[1] = {1.0f}; + ObservableArray src(arr, 1); + auto& f = src.HighPass(0.5); + Sink s; + f.Subscribe(s); + TEST_ASSERT_FLOAT_WITHIN(0.0001f, 0.5f, s.values[0]); +} + +void test_pass_stop_band(void) +{ + float arr[1] = {1.0f}; + ObservableArray src(arr, 1); + auto& pb = src.PassBand(0.1, 0.9); + Sink sp; + pb.Subscribe(sp); + TEST_ASSERT_FLOAT_WITHIN(0.0001f, 0.8f, sp.values[0]); // high(0.9) - low(0.1) + + auto& sb = src.StopBand(0.1, 0.9); + Sink ss; + sb.Subscribe(ss); + TEST_ASSERT_FLOAT_WITHIN(0.0001f, 0.2f, ss.values[0]); // 1.0 - 0.8 +} + +void test_is_equal(void) +{ + int arr[4] = {1, 2, 3, 0}; + ObservableArray src(arr, 4); + auto& f = src.IsEqual(2); + Sink s; + f.Subscribe(s); + TEST_ASSERT_EQUAL_INT(1, s.values.size()); + TEST_ASSERT_EQUAL_INT(2, s.values[0]); +} + +void test_is_less(void) +{ + int arr[4] = {1, 2, 3, 0}; + ObservableArray src(arr, 4); + auto& f = src.IsLess(3); + Sink s; + f.Subscribe(s); + TEST_ASSERT_EQUAL_INT(3, s.values.size()); // 1, 2, 0 are < 3 +} + +void test_is_zero(void) +{ + int arr[4] = {1, 2, 3, 0}; + ObservableArray src(arr, 4); + auto& f = src.IsZero(); + Sink s; + f.Subscribe(s); + TEST_ASSERT_EQUAL_INT(1, s.values.size()); + TEST_ASSERT_EQUAL_INT(0, s.values[0]); +} + +void test_debounce(void) +{ + ObservableProperty src; + auto& f = src.DebounceMillis(10); + Sink s; + f.Subscribe(s); + g_millis = 0; src = 1; // elapsed 0 -> dropped + g_millis = 20; src = 2; // elapsed 20 -> emitted + g_millis = 25; src = 3; // elapsed 5 -> dropped + g_millis = 35; src = 4; // elapsed 15 -> emitted + TEST_ASSERT_EQUAL_INT(2, s.values.size()); + TEST_ASSERT_EQUAL_INT(2, s.values[0]); + TEST_ASSERT_EQUAL_INT(4, s.values[1]); +} + +void test_window(void) +{ + ObservableProperty src; + auto& f = src.WindowMillis(100); + Sink s; + f.Subscribe(s); + g_millis = 0; src = 1; // opens window, elapsed 0 -> emitted + g_millis = 50; src = 2; // elapsed 50 -> emitted + g_millis = 200; src = 3; // elapsed 200 -> dropped + TEST_ASSERT_EQUAL_INT(2, s.values.size()); +} diff --git a/test/test_main.cpp b/test/test_main.cpp new file mode 100644 index 0000000..0ac623e --- /dev/null +++ b/test/test_main.cpp @@ -0,0 +1,215 @@ +// Single Unity test runner for the ReactiveArduino unit tests. +// +// PlatformIO links every *.cpp file in `test/` into a single test binary, so +// there must be exactly one main(), setUp() and tearDown(). Test functions live +// in the per-topic test_*.cpp files and are declared/run here. +#include +#include "ReactiveArduinoLib.h" +using namespace Reactive; +#include "TestHelpers.h" + +void setUp(void) { resetMocks(); } +void tearDown(void) {} + +// ---- test_aggregates.cpp ---- +void test_count(void); +void test_countdown(void); +void test_sum(void); +void test_min(void); +void test_max(void); +void test_average_float(void); +void test_rms_float(void); +void test_any(void); +void test_all_true(void); +void test_all_detects_failure(void); +void test_none(void); +void test_none_detects_match(void); + +// ---- test_filters.cpp ---- +void test_median3_bruteforce(void); +void test_median5_bruteforce(void); +void test_moving_average(void); +void test_moving_rms(void); +void test_on_rising(void); +void test_on_falling(void); +void test_low_pass(void); +void test_high_pass(void); +void test_pass_stop_band(void); +void test_is_equal(void); +void test_is_less(void); +void test_is_zero(void); +void test_debounce(void); +void test_window(void); + +// ---- test_observables.cpp ---- +void test_range_ascending(void); +void test_range_descending(void); +void test_range_defer(void); +void test_array(void); +void test_array_defer(void); +void test_property(void); +void test_manual_defer(void); +void test_timer_one_shot(void); +void test_timer_rearm(void); +void test_interval_periodic(void); +void test_analog_input(void); +void test_digital_input(void); + +// ---- test_observers.cpp ---- +void test_do(void); +void test_do_nothing(void); +void test_finally(void); +void test_do_and_finally(void); +void test_to_property(void); +void test_to_array(void); +void test_to_circular_buffer(void); +void test_digital_output(void); +void test_analog_output(void); +void test_serial_output(void); + +// ---- test_operators.cpp ---- +void test_where(void); +void test_distinct(void); +void test_first(void); +void test_last(void); +void test_skip(void); +void test_take(void); +void test_take_at(void); +void test_take_first(void); +void test_take_last(void); +void test_take_until(void); +void test_take_while(void); +void test_skip_until(void); +void test_skip_while(void); +void test_batch(void); +void test_foreach(void); +void test_if(void); +void test_timeout_millis(void); +void test_repeat(void); +void test_do_reset(void); +void test_not_reset(void); + +// ---- test_transformations.cpp ---- +void test_select(void); +void test_map(void); +void test_cast(void); +void test_reduce(void); +void test_limit(void); +void test_limit_upper(void); +void test_limit_lower(void); +void test_scale(void); +void test_abs(void); +void test_adc_to_voltage(void); +void test_toggle(void); +void test_threshold(void); +void test_elapsed_millis(void); +void test_timestamp_millis(void); +void test_frequency(void); +void test_to_bool(void); +void test_string_buffer(void); +void test_split(void); +void test_join(void); +void test_parse_int(void); +void test_parse_float(void); + +int main(void) +{ + UNITY_BEGIN(); + + RUN_TEST(test_count); + RUN_TEST(test_countdown); + RUN_TEST(test_sum); + RUN_TEST(test_min); + RUN_TEST(test_max); + RUN_TEST(test_average_float); + RUN_TEST(test_rms_float); + RUN_TEST(test_any); + RUN_TEST(test_all_true); + RUN_TEST(test_all_detects_failure); + RUN_TEST(test_none); + RUN_TEST(test_none_detects_match); + + RUN_TEST(test_median3_bruteforce); + RUN_TEST(test_median5_bruteforce); + RUN_TEST(test_moving_average); + RUN_TEST(test_moving_rms); + RUN_TEST(test_on_rising); + RUN_TEST(test_on_falling); + RUN_TEST(test_low_pass); + RUN_TEST(test_high_pass); + RUN_TEST(test_pass_stop_band); + RUN_TEST(test_is_equal); + RUN_TEST(test_is_less); + RUN_TEST(test_is_zero); + RUN_TEST(test_debounce); + RUN_TEST(test_window); + + RUN_TEST(test_range_ascending); + RUN_TEST(test_range_descending); + RUN_TEST(test_range_defer); + RUN_TEST(test_array); + RUN_TEST(test_array_defer); + RUN_TEST(test_property); + RUN_TEST(test_manual_defer); + RUN_TEST(test_timer_one_shot); + RUN_TEST(test_timer_rearm); + RUN_TEST(test_interval_periodic); + RUN_TEST(test_analog_input); + RUN_TEST(test_digital_input); + + RUN_TEST(test_do); + RUN_TEST(test_do_nothing); + RUN_TEST(test_finally); + RUN_TEST(test_do_and_finally); + RUN_TEST(test_to_property); + RUN_TEST(test_to_array); + RUN_TEST(test_to_circular_buffer); + RUN_TEST(test_digital_output); + RUN_TEST(test_analog_output); + RUN_TEST(test_serial_output); + + RUN_TEST(test_where); + RUN_TEST(test_distinct); + RUN_TEST(test_first); + RUN_TEST(test_last); + RUN_TEST(test_skip); + RUN_TEST(test_take); + RUN_TEST(test_take_at); + RUN_TEST(test_take_first); + RUN_TEST(test_take_last); + RUN_TEST(test_take_until); + RUN_TEST(test_take_while); + RUN_TEST(test_skip_until); + RUN_TEST(test_skip_while); + RUN_TEST(test_batch); + RUN_TEST(test_foreach); + RUN_TEST(test_if); + RUN_TEST(test_timeout_millis); + RUN_TEST(test_repeat); + RUN_TEST(test_do_reset); + RUN_TEST(test_not_reset); + + RUN_TEST(test_select); + RUN_TEST(test_map); + RUN_TEST(test_cast); + RUN_TEST(test_reduce); + RUN_TEST(test_limit); + RUN_TEST(test_limit_upper); + RUN_TEST(test_limit_lower); + RUN_TEST(test_scale); + RUN_TEST(test_abs); + RUN_TEST(test_adc_to_voltage); + RUN_TEST(test_toggle); + RUN_TEST(test_threshold); + RUN_TEST(test_elapsed_millis); + RUN_TEST(test_timestamp_millis); + RUN_TEST(test_frequency); + RUN_TEST(test_to_bool); + RUN_TEST(test_string_buffer); + RUN_TEST(test_split); + RUN_TEST(test_join); + RUN_TEST(test_parse_int); + RUN_TEST(test_parse_float); + + return UNITY_END(); +} diff --git a/test/test_observables.cpp b/test/test_observables.cpp new file mode 100644 index 0000000..e85ca66 --- /dev/null +++ b/test/test_observables.cpp @@ -0,0 +1,137 @@ +#include +#include "ReactiveArduinoLib.h" +using namespace Reactive; +#include "TestHelpers.h" + +void test_range_ascending(void) +{ + ObservableRange src(1, 5); + Sink s; + src.Subscribe(s); + TEST_ASSERT_EQUAL_INT(5, s.values.size()); + TEST_ASSERT_EQUAL_INT(1, s.values[0]); + TEST_ASSERT_EQUAL_INT(5, s.values.back()); +} + +void test_range_descending(void) +{ + ObservableRange src(5, 1, -1); + Sink s; + src.Subscribe(s); + TEST_ASSERT_EQUAL_INT(5, s.values.size()); + TEST_ASSERT_EQUAL_INT(1, s.values.back()); +} + +void test_range_defer(void) +{ + ObservableRangeDefer src(1, 3); + Sink s; + src.Subscribe(s); + src.Next(); + src.Next(); + src.Next(); + src.Next(); // past the end + TEST_ASSERT_EQUAL_INT(3, s.values.size()); + TEST_ASSERT_EQUAL_INT(1, s.completeCount); +} + +void test_array(void) +{ + int arr[3] = {10, 20, 30}; + ObservableArray src(arr, 3); + Sink s; + src.Subscribe(s); + TEST_ASSERT_EQUAL_INT(3, s.values.size()); + TEST_ASSERT_EQUAL_INT(30, s.values.back()); +} + +void test_array_defer(void) +{ + int arr[3] = {1, 2, 3}; + ObservableArrayDefer src(arr, 3); + Sink s; + src.Subscribe(s); + src.Next(); + src.Next(); + src.Next(); + src.Next(); + TEST_ASSERT_EQUAL_INT(3, s.values.size()); +} + +void test_property(void) +{ + ObservableProperty src; + Sink s; + src.Subscribe(s); + src = 7; + src = 8; + TEST_ASSERT_EQUAL_INT(2, s.values.size()); + src.Finish(); + src = 9; // ignored after completion + TEST_ASSERT_EQUAL_INT(2, s.values.size()); +} + +void test_manual_defer(void) +{ + ObservableManualDefer src; + Sink s; + src.Subscribe(s); + src.Next(); + src.Next(); + TEST_ASSERT_EQUAL_INT(2, s.values.size()); +} + +void test_timer_one_shot(void) +{ + ObservableTimerMillis t(1000); + Sink s; + t.Subscribe(s); + g_millis = 500; t.Update(); // not yet + g_millis = 1000; t.Update(); // fire + g_millis = 2000; t.Update(); // already expired + TEST_ASSERT_EQUAL_INT(1, s.values.size()); + TEST_ASSERT_TRUE(s.values[0] == 1000UL); +} + +void test_timer_rearm(void) +{ + ObservableTimerMillis t(1000); + Sink s; + t.Subscribe(s); + g_millis = 1000; t.Update(); // fire + t.Reset(); // re-arm (startTime = 1000) + g_millis = 1500; t.Update(); // not yet + g_millis = 2000; t.Update(); // fire again + TEST_ASSERT_EQUAL_INT(2, s.values.size()); +} + +void test_interval_periodic(void) +{ + ObservableIntervalMillis iv(1000); + Sink s; + iv.Subscribe(s); + g_millis = 1000; iv.Update(); // fire + g_millis = 1500; iv.Update(); // not yet + g_millis = 2000; iv.Update(); // fire + TEST_ASSERT_EQUAL_INT(2, s.values.size()); +} + +void test_analog_input(void) +{ + ObservableAnalogInput src(0); + Sink s; + src.Subscribe(s); + g_analogRead = 42; + src.Next(); + TEST_ASSERT_EQUAL_INT(42, s.values.back()); +} + +void test_digital_input(void) +{ + ObservableDigitalInput src(1); + Sink s; + src.Subscribe(s); + g_digitalRead = 1; + src.Next(); + TEST_ASSERT_EQUAL_INT(1, s.values.back()); +} diff --git a/test/test_observers.cpp b/test/test_observers.cpp new file mode 100644 index 0000000..c6f264e --- /dev/null +++ b/test/test_observers.cpp @@ -0,0 +1,106 @@ +#include +#include "ReactiveArduinoLib.h" +using namespace Reactive; +#include "TestHelpers.h" + +static int g_actionCount = 0; +static void countAction(int) { g_actionCount++; } +static int g_cbCount = 0; +static void countCb() { g_cbCount++; } + +void test_do(void) +{ + g_actionCount = 0; + g_cbCount = 0; + int arr[3] = {1, 2, 3}; + ObservableArray src(arr, 3); + src.Do(countAction); + TEST_ASSERT_EQUAL_INT(3, g_actionCount); +} + +void test_do_nothing(void) +{ + int arr[2] = {1, 2}; + ObservableArray src(arr, 2); + src.DoNothing(); // must not crash + TEST_ASSERT_TRUE(true); +} + +void test_finally(void) +{ + g_actionCount = 0; + g_cbCount = 0; + int arr[2] = {1, 2}; + ObservableArray src(arr, 2); + src.Finally(countCb); + TEST_ASSERT_EQUAL_INT(1, g_cbCount); +} + +void test_do_and_finally(void) +{ + g_actionCount = 0; + g_cbCount = 0; + int arr[2] = {1, 2}; + ObservableArray src(arr, 2); + src.DoAndFinally(countAction, countCb); + TEST_ASSERT_EQUAL_INT(2, g_actionCount); + TEST_ASSERT_EQUAL_INT(1, g_cbCount); +} + +void test_to_property(void) +{ + int arr[2] = {10, 20}; + ObservableArray src(arr, 2); + int out = 0; + src.ToProperty(out); + TEST_ASSERT_EQUAL_INT(20, out); +} + +void test_to_array(void) +{ + int arr[3] = {1, 2, 3}; + ObservableArray src(arr, 3); + int out[3] = {0, 0, 0}; + auto& a = src.ToArray(out, 3); + TEST_ASSERT_EQUAL_INT(3, out[2]); + TEST_ASSERT_EQUAL_INT(3, a.GetIndex()); +} + +void test_to_circular_buffer(void) +{ + int arr[5] = {1, 2, 3, 4, 5}; + ObservableArray src(arr, 5); + int out[3] = {0, 0, 0}; + src.ToCircularBuffer(out, 3); + // circular overwrite of a 3-slot buffer with 5 values -> {4, 5, 3} + TEST_ASSERT_EQUAL_INT(4, out[0]); + TEST_ASSERT_EQUAL_INT(5, out[1]); + TEST_ASSERT_EQUAL_INT(3, out[2]); +} + +void test_digital_output(void) +{ + int arr[1] = {1}; + ObservableArray src(arr, 1); + src.ToDigitalOutput(13); + TEST_ASSERT_EQUAL_INT(OUTPUT, (int)g_pinModeMode); + TEST_ASSERT_EQUAL_INT(13, (int)g_digitalPin); + TEST_ASSERT_EQUAL_INT(1, (int)g_digitalValue); +} + +void test_analog_output(void) +{ + int arr[1] = {128}; + ObservableArray src(arr, 1); + src.ToAnalogOutput(9); + TEST_ASSERT_EQUAL_INT(9, (int)g_analogPin); + TEST_ASSERT_EQUAL_INT(128, g_analogValue); +} + +void test_serial_output(void) +{ + int arr[1] = {42}; + ObservableArray src(arr, 1); + src.ToSerial(); + TEST_ASSERT_EQUAL_INT(1, Serial.printCount); +} diff --git a/test/test_operators.cpp b/test/test_operators.cpp new file mode 100644 index 0000000..d5e1f47 --- /dev/null +++ b/test/test_operators.cpp @@ -0,0 +1,239 @@ +#include +#include "ReactiveArduinoLib.h" +using namespace Reactive; +#include "TestHelpers.h" + +static bool isEven(int v) { return v % 2 == 0; } +static bool lessThan3(int v) { return v < 3; } +static bool greaterThan3(int v) { return v > 3; } +static int g_actionCount = 0; +static void countAction(int) { g_actionCount++; } +static int g_cbCount = 0; +static void countCb() { g_cbCount++; } + +void test_where(void) +{ + int arr[6] = {1, 2, 3, 4, 5, 6}; + ObservableArray src(arr, 6); + auto& op = src.Where(isEven); + Sink s; + op.Subscribe(s); + TEST_ASSERT_EQUAL_INT(3, s.values.size()); + TEST_ASSERT_EQUAL_INT(6, s.values.back()); +} + +void test_distinct(void) +{ + int arr[6] = {1, 1, 2, 2, 3, 3}; + ObservableArray src(arr, 6); + auto& op = src.Distinct(); + Sink s; + op.Subscribe(s); + TEST_ASSERT_EQUAL_INT(3, s.values.size()); +} + +void test_first(void) +{ + int arr[5] = {1, 2, 3, 4, 5}; + ObservableArray src(arr, 5); + auto& op = src.First(); + Sink s; + op.Subscribe(s); + TEST_ASSERT_EQUAL_INT(1, s.values.size()); + TEST_ASSERT_EQUAL_INT(1, s.values[0]); + TEST_ASSERT_EQUAL_INT(1, s.completeCount); +} + +void test_last(void) +{ + int arr[5] = {1, 2, 3, 4, 5}; + ObservableArray src(arr, 5); + auto& op = src.Last(); + Sink s; + op.Subscribe(s); + TEST_ASSERT_EQUAL_INT(1, s.values.size()); + TEST_ASSERT_EQUAL_INT(5, s.values[0]); +} + +void test_skip(void) +{ + int arr[6] = {1, 2, 3, 4, 5, 6}; + ObservableArray src(arr, 6); + auto& op = src.Skip(2); + Sink s; + op.Subscribe(s); + TEST_ASSERT_EQUAL_INT(4, s.values.size()); + TEST_ASSERT_EQUAL_INT(3, s.values[0]); +} + +void test_take(void) +{ + int arr[6] = {1, 2, 3, 4, 5, 6}; + ObservableArray src(arr, 6); + auto& op = src.Take(3); + Sink s; + op.Subscribe(s); + TEST_ASSERT_EQUAL_INT(3, s.values.size()); + TEST_ASSERT_EQUAL_INT(3, s.values.back()); +} + +void test_take_at(void) +{ + int arr[5] = {1, 2, 3, 4, 5}; + ObservableArray src(arr, 5); + auto& op = src.TakeAt(2); + Sink s; + op.Subscribe(s); + TEST_ASSERT_EQUAL_INT(1, s.values.size()); + TEST_ASSERT_EQUAL_INT(3, s.values[0]); +} + +void test_take_first(void) +{ + int arr[5] = {1, 2, 3, 4, 5}; + ObservableArray src(arr, 5); + auto& op = src.TakeFirst(); + Sink s; + op.Subscribe(s); + TEST_ASSERT_EQUAL_INT(1, s.values.size()); + TEST_ASSERT_EQUAL_INT(1, s.values[0]); +} + +void test_take_last(void) +{ + int arr[5] = {1, 2, 3, 4, 5}; + ObservableArray src(arr, 5); + auto& op = src.TakeLast(); + Sink s; + op.Subscribe(s); + TEST_ASSERT_EQUAL_INT(1, s.values.size()); + TEST_ASSERT_EQUAL_INT(5, s.values[0]); +} + +void test_take_until(void) +{ + int arr[6] = {1, 2, 3, 4, 5, 6}; + ObservableArray src(arr, 6); + auto& op = src.TakeUntil(greaterThan3); + Sink s; + op.Subscribe(s); + TEST_ASSERT_EQUAL_INT(3, s.values.size()); // 1, 2, 3 +} + +void test_take_while(void) +{ + int arr[6] = {1, 2, 3, 4, 5, 6}; + ObservableArray src(arr, 6); + auto& op = src.TakeWhile(lessThan3); + Sink s; + op.Subscribe(s); + TEST_ASSERT_EQUAL_INT(2, s.values.size()); // 1, 2 +} + +void test_skip_until(void) +{ + int arr[6] = {1, 2, 3, 4, 5, 6}; + ObservableArray src(arr, 6); + auto& op = src.SkipUntil(greaterThan3); + Sink s; + op.Subscribe(s); + TEST_ASSERT_EQUAL_INT(3, s.values.size()); // 4, 5, 6 + TEST_ASSERT_EQUAL_INT(4, s.values[0]); +} + +void test_skip_while(void) +{ + int arr[6] = {1, 2, 3, 4, 5, 6}; + ObservableArray src(arr, 6); + auto& op = src.SkipWhile(lessThan3); + Sink s; + op.Subscribe(s); + TEST_ASSERT_EQUAL_INT(4, s.values.size()); // 3, 4, 5, 6 + TEST_ASSERT_EQUAL_INT(3, s.values[0]); +} + +void test_batch(void) +{ + int arr[6] = {1, 2, 3, 4, 5, 6}; + ObservableArray src(arr, 6); + auto& op = src.Batch(3); + Sink s; + op.Subscribe(s); + TEST_ASSERT_EQUAL_INT(6, s.values.size()); // no values dropped at boundaries +} + +void test_foreach(void) +{ + g_actionCount = 0; + g_cbCount = 0; + int arr[3] = {1, 2, 3}; + ObservableArray src(arr, 3); + auto& op = src.ForEach(countAction); + Sink s; + op.Subscribe(s); + TEST_ASSERT_EQUAL_INT(3, g_actionCount); + TEST_ASSERT_EQUAL_INT(3, s.values.size()); +} + +void test_if(void) +{ + g_actionCount = 0; + g_cbCount = 0; + int arr[4] = {1, 2, 3, 4}; + ObservableArray src(arr, 4); + auto& op = src.If(isEven, countAction); + Sink s; + op.Subscribe(s); + TEST_ASSERT_EQUAL_INT(2, g_actionCount); // only for 2 and 4 + TEST_ASSERT_EQUAL_INT(4, s.values.size()); +} + +void test_timeout_millis(void) +{ + g_actionCount = 0; + g_cbCount = 0; + ObservableProperty src; + auto& op = src.TimeoutMillis(1000, countCb); + Sink s; + op.Subscribe(s); + src = 1; // resets the timer + g_millis = 500; op.Update(); // not expired + TEST_ASSERT_EQUAL_INT(0, g_cbCount); + g_millis = 1500; op.Update(); // expired + TEST_ASSERT_EQUAL_INT(1, g_cbCount); + TEST_ASSERT_EQUAL_INT(1, s.completeCount); +} + +void test_repeat(void) +{ + int arr[2] = {1, 2}; + ObservableArray src(arr, 2); + auto& op = src.Repeat(2); + Sink s; + op.Subscribe(s); + TEST_ASSERT_EQUAL_INT(4, s.values.size()); // 1, 2, 1, 2 + TEST_ASSERT_EQUAL_INT(1, s.values[2]); + TEST_ASSERT_EQUAL_INT(2, s.values[3]); +} + +void test_do_reset(void) +{ + int arr[2] = {1, 2}; + ObservableArray src(arr, 2); + auto& op = src.DoReset(); + Sink s; + op.Subscribe(s); + TEST_ASSERT_EQUAL_INT(2, s.values.size()); + op.Reset(); // resets the parent -> re-runs + TEST_ASSERT_EQUAL_INT(4, s.values.size()); +} + +void test_not_reset(void) +{ + ObservableManualDefer src; + auto& op = src.NotReset(); + Sink s; + op.Subscribe(s); + op.Reset(); // completes children without resetting the parent + TEST_ASSERT_EQUAL_INT(1, s.completeCount); +} diff --git a/test/test_transformations.cpp b/test/test_transformations.cpp new file mode 100644 index 0000000..09de901 --- /dev/null +++ b/test/test_transformations.cpp @@ -0,0 +1,244 @@ +#include +#include "ReactiveArduinoLib.h" +using namespace Reactive; +#include "TestHelpers.h" + +static int addOne(int v) { return v + 1; } +static float half(int v) { return (float)v / 2.0f; } +static int sumInt(int acc, int v) { return acc + v; } + +void test_select(void) +{ + int arr[3] = {1, 2, 3}; + ObservableArray src(arr, 3); + auto& t = src.Select(addOne); + Sink s; + t.Subscribe(s); + TEST_ASSERT_EQUAL_INT(3, s.values.size()); + TEST_ASSERT_EQUAL_INT(4, s.values.back()); +} + +void test_map(void) +{ + int arr[3] = {1, 2, 3}; + ObservableArray src(arr, 3); + auto& t = src.Map(half); + Sink s; + t.Subscribe(s); + TEST_ASSERT_EQUAL_INT(3, s.values.size()); + TEST_ASSERT_FLOAT_WITHIN(0.001f, 1.5f, s.values.back()); +} + +void test_cast(void) +{ + int arr[3] = {1, 2, 3}; + ObservableArray src(arr, 3); + auto& t = src.Cast(); + Sink s; + t.Subscribe(s); + TEST_ASSERT_FLOAT_WITHIN(0.001f, 3.0f, s.values.back()); +} + +void test_reduce(void) +{ + int arr[4] = {1, 2, 3, 4}; + ObservableArray src(arr, 4); + auto& t = src.Reduce(sumInt, 0); + Sink s; + t.Subscribe(s); + TEST_ASSERT_EQUAL_INT(10, s.values.back()); +} + +void test_limit(void) +{ + int arr[5] = {0, 1, 5, 9, 10}; + ObservableArray src(arr, 5); + auto& t = src.Limit(1, 9); + Sink s; + t.Subscribe(s); + TEST_ASSERT_EQUAL_INT(1, s.values[0]); + TEST_ASSERT_EQUAL_INT(9, s.values.back()); +} + +void test_limit_upper(void) +{ + int arr[4] = {0, 5, 10, 15}; + ObservableArray src(arr, 4); + auto& t = src.LimitUpper(10); + Sink s; + t.Subscribe(s); + TEST_ASSERT_EQUAL_INT(10, s.values.back()); +} + +void test_limit_lower(void) +{ + int arr[4] = {0, 5, 10, 15}; + ObservableArray src(arr, 4); + auto& t = src.LimitLower(5); + Sink s; + t.Subscribe(s); + TEST_ASSERT_EQUAL_INT(5, s.values[0]); +} + +void test_scale(void) +{ + float arr[3] = {0.0f, 5.0f, 10.0f}; + ObservableArray src(arr, 3); + auto& t = src.Scale(0.0f, 10.0f, 0.0f, 100.0f); + Sink s; + t.Subscribe(s); + TEST_ASSERT_FLOAT_WITHIN(0.001f, 50.0f, s.values[1]); + TEST_ASSERT_FLOAT_WITHIN(0.001f, 100.0f, s.values.back()); +} + +void test_abs(void) +{ + // TransformationAbs is not exposed through a fluent method, so it is + // exercised directly through its public Operator interface. + TransformationAbs t; + Sink s; + t._childObservers.Add(&s); + t.OnNext(-3); + t.OnNext(0); + t.OnNext(3); + TEST_ASSERT_EQUAL_INT(3, s.values.size()); + TEST_ASSERT_EQUAL_INT(3, s.values[0]); + TEST_ASSERT_EQUAL_INT(3, s.values[2]); +} + +void test_adc_to_voltage(void) +{ + int arr[2] = {0, 512}; + ObservableArray src(arr, 2); + auto& t = src.AdcToVoltage(); + Sink s; + t.Subscribe(s); + TEST_ASSERT_FLOAT_WITHIN(0.001f, (512.0f * 5.0f) / 1023.0f, s.values.back()); +} + +void test_toggle(void) +{ + int arr[3] = {1, 1, 1}; + ObservableArray src(arr, 3); + auto& t = src.Toggle(); + Sink s; + t.Subscribe(s); + TEST_ASSERT_EQUAL_INT(1, s.values[0]); // first toggle -> HIGH + TEST_ASSERT_EQUAL_INT(0, s.values[1]); + TEST_ASSERT_EQUAL_INT(1, s.values[2]); +} + +void test_threshold(void) +{ + int arr[4] = {0, 2, 8, 10}; + ObservableArray src(arr, 4); + auto& t = src.Threshold(5); + Sink s; + t.Subscribe(s); + TEST_ASSERT_EQUAL_INT(0, s.values[1]); // 2 is still LOW + TEST_ASSERT_EQUAL_INT(1, s.values.back()); // 8 and 10 are HIGH +} + +void test_elapsed_millis(void) +{ + ObservableProperty src; + auto& t = src.ElapsedMillis(); + Sink s; + t.Subscribe(s); + g_millis = 100; src = 1; + g_millis = 250; src = 2; + TEST_ASSERT_TRUE(s.values[0] == 100UL); + TEST_ASSERT_TRUE(s.values[1] == 150UL); +} + +void test_timestamp_millis(void) +{ + ObservableProperty src; + auto& t = src.Millis(); + Sink s; + t.Subscribe(s); + g_millis = 100; src = 1; + g_millis = 250; src = 2; + TEST_ASSERT_TRUE(s.values[0] == 100UL); + TEST_ASSERT_TRUE(s.values[1] == 250UL); +} + +void test_frequency(void) +{ + ObservableProperty src; + auto& t = src.Frequency(); + Sink s; + t.Subscribe(s); + g_millis = 100; src = 1; // 1000 / 100 = 10 Hz + g_millis = 200; src = 2; // 1000 / 100 = 10 Hz + TEST_ASSERT_FLOAT_WITHIN(0.01f, 10.0f, s.values[0]); + TEST_ASSERT_FLOAT_WITHIN(0.01f, 10.0f, s.values[1]); +} + +void test_to_bool(void) +{ + int arr[3] = {0, 1, 2}; + ObservableArray src(arr, 3); + auto& t = src.ToBool(); + Sink s; + t.Subscribe(s); + TEST_ASSERT_FALSE(s.values[0]); + TEST_ASSERT_TRUE(s.values[1]); + TEST_ASSERT_TRUE(s.values[2]); +} + +void test_string_buffer(void) +{ + String sarr[2] = {String("ab"), String("cd")}; + ObservableArray src(sarr, 2); + auto& t = src.StringBuffer(); + Sink s; + t.Subscribe(s); + TEST_ASSERT_EQUAL_INT(2, s.values.size()); + TEST_ASSERT_TRUE(strcmp(s.values[1].c_str(), "abcd") == 0); +} + +void test_split(void) +{ + String sarr[1] = {String("a;b;c")}; + ObservableArray src(sarr, 1); + auto& t = src.Split(';'); + Sink s; + t.Subscribe(s); + TEST_ASSERT_EQUAL_INT(3, s.values.size()); + TEST_ASSERT_TRUE(strcmp(s.values[0].c_str(), "a") == 0); + TEST_ASSERT_TRUE(strcmp(s.values[2].c_str(), "c") == 0); +} + +void test_join(void) +{ + String sarr[3] = {String("a"), String("b"), String("c")}; + ObservableArray src(sarr, 3); + auto& t = src.Join('-'); + Sink s; + t.Subscribe(s); + TEST_ASSERT_EQUAL_INT(3, s.values.size()); + TEST_ASSERT_TRUE(strcmp(s.values[2].c_str(), "a-b-c") == 0); +} + +void test_parse_int(void) +{ + String sarr[2] = {String("123"), String("45")}; + ObservableArray src(sarr, 2); + auto& t = src.ParseInt(); + Sink s; + t.Subscribe(s); + TEST_ASSERT_EQUAL_INT(123, s.values[0]); + TEST_ASSERT_EQUAL_INT(45, s.values[1]); +} + +void test_parse_float(void) +{ + String sarr[2] = {String("3.5"), String("2.25")}; + ObservableArray src(sarr, 2); + auto& t = src.ParseFloat(); + Sink s; + t.Subscribe(s); + TEST_ASSERT_FLOAT_WITHIN(0.001f, 3.5f, s.values[0]); + TEST_ASSERT_FLOAT_WITHIN(0.001f, 2.25f, s.values[1]); +}