RxRust: a zero cost Rust implementation of Reactive Extensions
Usage
Add this to your Cargo.toml:
[dependencies]
rxrust = "0.3.0";
Example
use ;
let mut numbers = from_iter!;
// crate a even stream by filter
let even = numbers.fork.filter;
// crate an odd stream by filter
let odd = numbers.fork.filter;
// merge odd and even stream again
even.merge.subscribe;
// "0 1 2 3 4 5 6 7 8 9" will be printed.
Fork Stream
In rxrust
almost all extensions consume the upstream. So when you try to subscribe a stream twice, the compiler will complain.
# use *;
let o = from_iter!;
o.subscribe;
o.subscribe;
In this case, we can use Fork
to fork a stream. In general, Fork
has same mean with Clone
, this will not change a cold stream to hot stream. If you want convert a stream from unicast to multicast, from cold to hot use Publish
and RefCount
.
# use *;
# use Fork;
let o = from_iter!;
o.fork.subscribe;
o.fork.subscribe;
Scheduler
use *;
use ;
from_iter!
.subscribe_on
.map
.observe_on
.subscribe;
Converts from a Future
just use observable::from_future!
to convert a Future
to an observable sequence.
use *;
use future;
from_future!
.subscribe;
// because all future in rxrust are execute async, so we wait a second to see
// the print, no need to use this line in your code.
sleep;
A from_future_with_err!
macro also provided to propagating error from Future
.
All contributions are welcome
We are looking for contributors! Feel free to open issues for asking questions, suggesting features or other things!
Help and contributions can be any of the following:
- use the project and report issues to the project issues page
- documentation and README enhancement (VERY important)
- continuous improvement in a ci Pipeline
- implement any unimplemented operator, remember to create a pull request before you start your code, so other people know you are work on it.