Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 15 additions & 0 deletions .github/workflows/main.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
143 changes: 141 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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<int>();
obs >> ToSerial<int>();
// ...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.
Expand Down
Binary file removed ReactiveArduino Legend.png
Binary file not shown.
19 changes: 19 additions & 0 deletions platformio.ini
Original file line number Diff line number Diff line change
@@ -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
2 changes: 1 addition & 1 deletion src/Aggregates/AggregateAll.h
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ AggregateAll<T>::AggregateAll(ReactivePredicate<T> condition)
template <typename T>
void AggregateAll<T>::OnNext(T value)
{
if (_state && _condition(value)) _state = false;
if (!_condition(value)) _state = false;

this->_childObservers.OnNext(_state);
}
Expand Down
2 changes: 1 addition & 1 deletion src/Aggregates/AggregateAverage.h
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ void AggregateAverage<T>::OnNext(T value)
{
_sum += value;
_count++;
this->_childObservers.OnNext(_sum / _count);
this->_childObservers.OnNext(_sum / static_cast<T>(_count));
}

#endif
2 changes: 1 addition & 1 deletion src/Aggregates/AggregateCount.h
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ class AggregateCount : public Operator<T, int>
void OnNext(T value) override;

private:
int _count = false;
int _count = 0;
};

template <typename T>
Expand Down
6 changes: 5 additions & 1 deletion src/Aggregates/AggregateCountdown.h
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,8 @@ class AggregateCountdown : public Operator<T, int>
void OnNext(T value) override;

private:
int _count = false;
int _count = 0;
bool _completed = false;
};

template <typename T>
Expand All @@ -31,11 +32,14 @@ AggregateCountdown<T>::AggregateCountdown(int count)
template <typename T>
void AggregateCountdown<T>::OnNext(T value)
{
if (_completed) return;

_count--;
this->_childObservers.OnNext(_count);

if (_count <= 0)
{
_completed = true;
this->_childObservers.OnComplete();
}
}
Expand Down
2 changes: 1 addition & 1 deletion src/Aggregates/AggregateNone.h
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ AggregateNone<T>::AggregateNone(ReactivePredicate<T> condition)
template <typename T>
void AggregateNone<T>::OnNext(T value)
{
if (!_condition(value)) _state = false;
if (_condition(value)) _state = false;

this->_childObservers.OnNext(_state);
}
Expand Down
2 changes: 1 addition & 1 deletion src/Aggregates/AggregateRMS.h
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ void AggregateRMS<T>::OnNext(T value)
{
_sumSqr += value * value;
_count++;
this->_childObservers.OnNext(sqrt(_sumSqr / _count));
this->_childObservers.OnNext(static_cast<T>(sqrt(static_cast<double>(_sumSqr) / static_cast<double>(_count))));
}

#endif
2 changes: 0 additions & 2 deletions src/Filters/FilterWindowMicros.h
Original file line number Diff line number Diff line change
Expand Up @@ -32,9 +32,7 @@ void FilterWindowMicros<T>::OnNext(T value)
}

if (_started && static_cast<unsigned long>(micros() - _lastTrigger) <= _interval)
{
this->_childObservers.OnNext(value);
}
}

#endif
2 changes: 0 additions & 2 deletions src/Filters/FilterWindowMillis.h
Original file line number Diff line number Diff line change
Expand Up @@ -39,9 +39,7 @@ void FilterWindowMillis<T>::OnNext(T value)
}

if (_started && static_cast<unsigned long>(millis() - _lastTrigger) <= _interval)
{
this->_childObservers.OnNext(value);
}
}

#endif
5 changes: 3 additions & 2 deletions src/Observables/ObservableIntervalMicros.h
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,6 @@ class ObservableIntervalMicros : public Observable<unsigned long>

private:
bool _isActive;
bool _isExpired;
unsigned long _startTime;
unsigned long _delay;
unsigned long _offset;
Expand Down Expand Up @@ -121,7 +120,9 @@ unsigned long ObservableIntervalMicros<T>::GetElapsedTime()
template <typename T>
unsigned long ObservableIntervalMicros<T>::GetRemainingTime()
{
return _interval - micros() + _startTime;
unsigned long elapsed = micros() - _startTime;
if (elapsed >= _interval) return 0;
return _interval - elapsed;
}

template <typename T>
Expand Down
5 changes: 3 additions & 2 deletions src/Observables/ObservableIntervalMillis.h
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,6 @@ class ObservableIntervalMillis : public Observable<unsigned long>

private:
bool _isActive;
bool _isExpired;
unsigned long _startTime;
unsigned long _delay;
unsigned long _offset;
Expand Down Expand Up @@ -121,7 +120,9 @@ unsigned long ObservableIntervalMillis<T>::GetElapsedTime()
template <typename T>
unsigned long ObservableIntervalMillis<T>::GetRemainingTime()
{
return _interval - millis() + _startTime;
unsigned long elapsed = millis() - _startTime;
if (elapsed >= _interval) return 0;
return _interval - elapsed;
}

template <typename T>
Expand Down
12 changes: 10 additions & 2 deletions src/Observables/ObservableRange.h
Original file line number Diff line number Diff line change
Expand Up @@ -54,8 +54,16 @@ void ObservableRange<T>::UnSubscribe(IObserver<T> &observer)
template<typename T>
void ObservableRange<T>::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();
}
Expand Down
5 changes: 3 additions & 2 deletions src/Observables/ObservableRangeDefer.h
Original file line number Diff line number Diff line change
Expand Up @@ -53,13 +53,14 @@ void ObservableRangeDefer<T>::UnSubscribe(IObserver<T> &observer)
template <typename T>
void ObservableRangeDefer<T>::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();
}

Expand Down
4 changes: 2 additions & 2 deletions src/Observables/ObservableSerialByte.h
Original file line number Diff line number Diff line change
Expand Up @@ -15,8 +15,8 @@ class ObservableSerial<byte> : public Observable<byte>
{
public:
ObservableSerial();
void Subscribe(IObserver<byte> &observer);
void UnSubscribe(IObserver<byte> &observer);
void Subscribe(IObserver<byte> &observer) override;
void UnSubscribe(IObserver<byte> &observer) override;
void Receive();

private:
Expand Down
2 changes: 1 addition & 1 deletion src/Observables/ObservableSerialDouble.h
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ class ObservableSerial<double> : public Observable<double>
private:
char _separator;

float _data = 0;
double _data = 0;
int _dataReal = 0;
int _dataDecimal = 0;
int _dataPow = 1;
Expand Down
2 changes: 0 additions & 2 deletions src/Observables/ObservableSerialString.h
Original file line number Diff line number Diff line change
Expand Up @@ -47,9 +47,7 @@ inline void ObservableSerial<String>::Receive()
{
const char newChar = Serial.read();
if (newChar != _separator)
{
_buffer.concat(newChar);
}
else
{
_childObservers.OnNext(_buffer);
Expand Down
Loading
Loading