diff --git a/.github/dependabot.yml b/.github/dependabot.yml new file mode 100644 index 0000000..0d08e26 --- /dev/null +++ b/.github/dependabot.yml @@ -0,0 +1,11 @@ +# To get started with Dependabot version updates, you'll need to specify which +# package ecosystems to update and where the package manifests are located. +# Please see the documentation for all configuration options: +# https://docs.github.com/code-security/dependabot/dependabot-version-updates/configuration-options-for-the-dependabot.yml-file + +version: 2 +updates: + - package-ecosystem: "github-actions" # See documentation for possible values + directory: "/" # Location of package manifests + schedule: + interval: "weekly" diff --git a/.github/workflows/main.yml b/.github/workflows/main.yml index 3d482c4..8a793c5 100644 --- a/.github/workflows/main.yml +++ b/.github/workflows/main.yml @@ -8,6 +8,20 @@ jobs: matrix: pio-env: ['esp12e', 'esp32dev'] example: ['Blink/Blink.cpp', 'Do/Do.cpp', 'FromArray/FromArray.cpp', 'FromProperty/FromProperty.cpp', 'FromSerialPort/FromSerialPort.cpp', 'FromString/FromString.cpp', 'GlobalDefined/GlobalDefined.cpp', 'Reduce/Reduce.cpp'] + steps: + - uses: actions/checkout@v4 + - name: Set up python + uses: actions/setup-python@v5 + with: + python-version: '3.x' + - name: Install PlatformIO + 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 @@ -17,5 +31,5 @@ jobs: architecture: 'x64' - name: Install PlatformIO run: python -m pip install platformio - - name: Build firmware - run: pio ci --lib="." --board ${{matrix.pio-env}} "examples/${{matrix.example}}" + - name: Run unit tests + run: pio test -e native diff --git a/keywords.txt b/keywords.txt deleted file mode 100644 index a4247cd..0000000 --- a/keywords.txt +++ /dev/null @@ -1,16 +0,0 @@ -####################################### -# Syntax Coloring Map ReactiveArduino -####################################### - -####################################### -# Datatypes (KEYWORD1) -####################################### -ReactiveArduino KEYWORD1 - -####################################### -# Methods and Functions (KEYWORD2) -####################################### - -####################################### -# Constants (LITERAL1) -####################################### diff --git a/library.json b/library.json new file mode 100644 index 0000000..1c10ec7 --- /dev/null +++ b/library.json @@ -0,0 +1,37 @@ +{ + "$schema": "https://raw.githubusercontent.com/platformio/platformio-core/develop/platformio/assets/schema/library.json", + "name": "ReactiveArduino", + "description": "ReactiveArduino implements observable-observer pattern on a processor like Arduino", + "keywords": [ + "Reactive", + "Arduino", + "Observer", + "Observable", + "Event" + ], + "repository": { + "type": "git", + "url": "https://github.com/rzeldent/Arduino-ReactiveArduino" + }, + "frameworks": [ + "arduino" + ], + "platforms": "*", + "build": { + "srcDir": "src/", + "includeDir": "include/" + }, + "version": "2.0.0", + "authors": [ + { + "name": "Luis Llamas", + "maintainer": true + }, + { + "name": "Rene Zeldenthuis", + "maintainer": true + } + + ], + "category": "Other" +} \ No newline at end of file diff --git a/library.properties b/library.properties deleted file mode 100644 index f36f703..0000000 --- a/library.properties +++ /dev/null @@ -1,9 +0,0 @@ -name=ReactiveArduino -version=2.0.0 -author=Luis Llamas -maintainer=Luis Llamas -sentence=ReactiveArduino implements observable-observer pattern on a processor like Arduino -paragraph=ReactiveArduino implements observable-observer pattern on a processor like Arduino -category=Other -url=https://github.com/luisllamasbinaburo/Arduino-ReactiveArduino -architectures=* 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/OperatorBufferCount.h b/src/Operators/OperatorBufferCount.h new file mode 100644 index 0000000..e42c095 --- /dev/null +++ b/src/Operators/OperatorBufferCount.h @@ -0,0 +1,50 @@ +/*************************************************** +Copyright (c) 2019 Luis Llamas +(www.luisllamas.es) + +Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions and limitations under the License + ****************************************************/ + +#ifndef _REACTIVEOPERATORBUFFERCOUNT_h +#define _REACTIVEOPERATORBUFFERCOUNT_h + +template +class OperatorBufferCount: public Operator +{ +public: + OperatorBufferCount(size_t N); + + void OnNext(T value) override; + + void Reset() override; +private: + std::vector _buffer; + size_t _num_elements = 0; +}; + +template +OperatorBufferCount::OperatorBufferCount(size_t N) +{ + _num_elements = N; +} + +template +void OperatorBufferCount::Reset() +{ + _buffer.clear(); +} + +template +void OperatorBufferCount::OnNext(T value) +{ + _buffer.push_back(value); + if (_buffer.size() >= _num_elements) + { + this->_childObservers.OnNext(_buffer.data()); + _buffer.clear(); + } +} + +#endif \ No newline at end of file diff --git a/src/Operators/OperatorDistinct.h b/src/Operators/OperatorDistinct.h index 268a186..58a01aa 100644 --- a/src/Operators/OperatorDistinct.h +++ b/src/Operators/OperatorDistinct.h @@ -19,7 +19,7 @@ class OperatorDistinct : public Operator void OnNext(T value) override; private: - T _last = T(); + std::set _seen; bool _any = false; }; @@ -31,10 +31,12 @@ OperatorDistinct::OperatorDistinct() template void OperatorDistinct::OnNext(T value) { - if (!_any || (_any && _last != value)) + if (!_any || (_any && _seen.find(value) == _seen.end())) + { + _seen.insert(value); this->_childObservers.OnNext(value); + } - _last = value; _any = true; } diff --git a/src/Operators/OperatorDistinctUntilChanged.h b/src/Operators/OperatorDistinctUntilChanged.h new file mode 100644 index 0000000..c6ef9df --- /dev/null +++ b/src/Operators/OperatorDistinctUntilChanged.h @@ -0,0 +1,43 @@ +/*************************************************** +Copyright (c) 2019 Luis Llamas +(www.luisllamas.es) + +Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions and limitations under the License + ****************************************************/ + +#ifndef _REACTIVEOPERATORDISTINCTUNTILCHANGED_h +#define _REACTIVEOPERATORDISTINCTUNTILCHANGED_h + +template +class OperatorDistinctUntilChanged : public Operator +{ +public: + OperatorDistinctUntilChanged(); + + void OnNext(T value) override; + +private: + T _last = T(); + bool _any = false; +}; + +template +OperatorDistinctUntilChanged::OperatorDistinctUntilChanged() +{ +} + +template +void OperatorDistinctUntilChanged::OnNext(T value) +{ + if (!_any || (_any && _last != value)) + { + _last = value; + this->_childObservers.OnNext(value); + } + + _any = true; +} + +#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/Operators/Operators.h b/src/Operators/Operators.h index 8b07170..0b97bcd 100644 --- a/src/Operators/Operators.h +++ b/src/Operators/Operators.h @@ -12,6 +12,7 @@ Unless required by applicable law or agreed to in writing, software distributed #include "OperatorWhere.h" #include "OperatorDistinct.h" +#include "OperatorDistinctUntilChanged.h" #include "OperatorLast.h" #include "OperatorFirst.h" @@ -19,6 +20,7 @@ Unless required by applicable law or agreed to in writing, software distributed #include "OperatorTake.h" #include "OperatorSkip.h" #include "OperatorBatch.h" +#include "OperatorBufferCount.h" #include "OperatorTakeAt.h" #include "OperatorTakeFirst.h" diff --git a/src/ReactiveArduinoCore.h b/src/ReactiveArduinoCore.h index 37e3ae6..d2eb48d 100644 --- a/src/ReactiveArduinoCore.h +++ b/src/ReactiveArduinoCore.h @@ -50,11 +50,13 @@ template class FilterIsZero; template class OperatorWhere; template class OperatorDistinct; +template class OperatorDistinctUntilChanged; template class OperatorLast; template class OperatorFirst; template class OperatorTake; template class OperatorSkip; template class OperatorBatch; +template class OperatorBufferCount; template class OperatorTakeAt; template class OperatorTakeFirst; template class OperatorTakeLast; @@ -154,6 +156,7 @@ class Observable : IObservable, IResetable // "Fluent" behavior OperatorWhere& Where(ReactivePredicate condition); OperatorDistinct& Distinct(); + OperatorDistinctUntilChanged& DistinctUntilChanged(); OperatorFirst& First(); OperatorLast& Last(); OperatorSkip& Skip(size_t num); @@ -166,10 +169,11 @@ class Observable : IObservable, IResetable OperatorTakeUntil& TakeUntil(ReactivePredicate condition); OperatorTakeWhile& TakeWhile(ReactivePredicate condition); OperatorBatch& Batch(size_t num); + OperatorBufferCount& BufferCount(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 +194,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 +241,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); @@ -266,6 +270,14 @@ auto Observable::Distinct() -> OperatorDistinct& return *newOp; } +template +auto Observable::DistinctUntilChanged() -> OperatorDistinctUntilChanged& +{ + auto newOp = new OperatorDistinctUntilChanged(); + Compound(*this, *newOp); + return *newOp; +} + template auto Observable::First() -> OperatorFirst& { @@ -368,6 +380,14 @@ auto Observable::Batch(size_t num) -> OperatorBatch& return *newOp; } +template +auto Observable::BufferCount(size_t num) -> OperatorBufferCount& +{ + auto newOp = new OperatorBufferCount(num); + Compound(*this, *newOp); + return *newOp; +} + template auto Observable::If(ReactivePredicate condition, ReactiveAction action) -> OperatorIf& { @@ -385,17 +405,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 +580,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 +588,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 +596,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 +604,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 +932,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..66cf67c 100644 --- a/src/ReactiveArduinoLib.h +++ b/src/ReactiveArduinoLib.h @@ -5,7 +5,7 @@ Copyright (c) 2019 Luis Llamas Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at http://www.apache.org/licenses/LICENSE-2.0 Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions and limitations under the License - ****************************************************/ +****************************************************/ #ifndef _REACTIVEARDUINOLIB_h #define _REACTIVEARDUINOLIB_h @@ -18,6 +18,9 @@ Unless required by applicable law or agreed to in writing, software distributed #include "WProgram.h" #endif +#include +#include + namespace Reactive { #include "ReactiveArduinoCore.h" @@ -111,7 +114,7 @@ namespace Reactive template auto FromSerial(char separator) -> ObservableSerial& { - return *(new ObservableSerial()); + return *(new ObservableSerial(separator)); } template <> @@ -223,6 +226,12 @@ namespace Reactive return *(new OperatorDistinct()); } + template + OperatorDistinctUntilChanged& DistinctUntilChanged() + { + return *(new OperatorDistinctUntilChanged()); + } + template OperatorFirst& First() { @@ -314,15 +323,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 +466,7 @@ namespace Reactive template TransformationThreshold& DoubleThreshold(T lowThreshold, T highThreshold) { - return *(new TransformationThreshold(lowThreshold, highThreshold)); + return *(new TransformationThreshold(lowThreshold, highThreshold, false)); } template @@ -473,7 +482,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]); +}