-
Notifications
You must be signed in to change notification settings - Fork 443
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Reactor Integration #706
Merged
Merged
Reactor Integration #706
Changes from all commits
Commits
Show all changes
21 commits
Select commit
Hold shift + click to select a range
2c7b427
Starting Reactor integration
Cotel b7ff184
Merge branch 'master' into cotel-reactor-integration
raulraja 7baa810
Adding some tests
Cotel 0ada9ac
Merge branch 'cotel-reactor-integration' of https://github.com/arrow-…
Cotel 84c2848
Adding instances for MonoK
Cotel 43cd3e3
Merge branch 'master' of https://github.com/arrow-kt/arrow into cotel…
Cotel bd1a5f6
Including arrow-effects-reactor in settings.gradle
Cotel 5b11f27
Updating FluxK
Cotel bac7a86
Updating MonoK
Cotel 0058221
Fixing tests
Cotel 178f6f8
Adding another test for FluxK
Cotel a5adc6e
Adding tests for MonoK
Cotel e06c40a
Merge branch 'master' into cotel-reactor-integration
pakoito 7218a0f
Merge branch 'master' into cotel-reactor-integration
pakoito db0d360
Removing duplicated laws for MonoK
Cotel a12cefd
Adding docs for reactor integration
Cotel 3778fe8
Merge remote-tracking branch 'origin/cotel-reactor-integration' into …
Cotel ef7e8f1
Including missing dependencies
Cotel 930a3f1
Fixing typo in package name
Cotel dfbc423
Adding reactor dependency line to the root README
Cotel bcf91c4
Merge branch 'master' into cotel-reactor-integration
pakoito File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Large diffs are not rendered by default.
Oops, something went wrong.
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
157 changes: 157 additions & 0 deletions
157
modules/docs/arrow-docs/docs/docs/integrations/reactor/README.md
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,157 @@ | ||
--- | ||
layout: docs | ||
title: Reactor | ||
permalink: /docs/integrations/reactor/ | ||
--- | ||
|
||
## Project Reactor | ||
|
||
Arrow aims to enhance the user experience when using Project Reactor. While providing other datatypes that are capable of handling effects, like IO, the style of programming encouraged by the library allows users to generify behavior for any existing abstractions. | ||
|
||
One of such abstractions is Project Reactor, a library that, like RxJava, offers reactive streams. | ||
|
||
```kotlin | ||
val flux = Flux.just(7, 4, 11 ,3) | ||
.map { it + 1 } | ||
.filter { it % 2 == 0 } | ||
.scan { acc, value -> acc + value } | ||
.collectList() | ||
.subscribeOn(Schedulers.parallel()) | ||
.block() | ||
//[8, 20, 24] | ||
``` | ||
|
||
### Integration with your existing Flux chains | ||
|
||
The largest quality of life improvement when using Flux streams in Arrow is the introduction of the [Monad Comprehension]({{ '/docs/patterns/monad_comprehensions' | relative_url }}). This library construct allows expressing asynchronous Flux sequences as synchronous code using binding/bind. | ||
|
||
#### Arrow Wrapper | ||
|
||
To wrap any existing Flux in its Arrow Wrapper counterpart you can use the extension function `k()`. | ||
|
||
```kotlin:ank | ||
import arrow.effects.* | ||
import reactor.core.publisher.* | ||
|
||
val flux = Flux.just(1, 2, 3, 4, 5).k() | ||
flux | ||
``` | ||
|
||
```kotlin:ank | ||
val mono = Mono.just(1).k() | ||
mono | ||
``` | ||
|
||
You can return to their regular forms using the function `value()`. | ||
|
||
```kotlin:ank | ||
flux.value() | ||
``` | ||
|
||
```kotlin:ank | ||
mono.value() | ||
``` | ||
|
||
### Observable comprehensions | ||
|
||
The library provides instances of [`MonadError`]({{ '/docs/typeclasses/monaderror' | relative_url }}) and [`MonadDefer`]({{ '/docs/effects/monaddefer' | relative_url }}). | ||
|
||
[`MonadDefer`]({{ '/docs/effects/async' | relative_url }}) allows you to generify over datatypes that can run asynchronous code. You can use it with `FluxK` or `MonoK`. | ||
|
||
```kotlin | ||
fun <F> getSongUrlAsync(MS: MonadDefer<F>) = | ||
MS { getSongUrl() } | ||
|
||
val songFlux: FluxKOf<Url> = getSongUrlAsync(FluxK.monadDefer()) | ||
val songMono: MonoKOf<Url> = getSongUrlAsync(MonoK.monadDefer()) | ||
``` | ||
|
||
[`MonadError`]({{ '/docs/typeclasses/monaderror' | relative_url }}) can be used to start a [Monad Comprehension]({{ '/docs/patterns/monad_comprehensions' | relative_url }}) using the method `bindingCatch`, with all its benefits. | ||
|
||
Let's take an example and convert it to a comprehension. We'll create an observable that loads a song from a remote location, and then reports the current play % every 100 milliseconds until the percentage reaches 100%: | ||
|
||
```kotlin | ||
getSongUrlAsync() | ||
.map { songUrl -> MediaPlayer.load(songUrl) } | ||
.flatMap { | ||
val totalTime = musicPlayer.getTotaltime() | ||
Flux.interval(Duration.ofMillis(100)) | ||
.flatMap { | ||
Flux.create { musicPlayer.getCurrentTime() } | ||
.map { tick -> (tick / totalTime * 100).toInt() } | ||
} | ||
.takeUntil { percent -> percent >= 100 } | ||
} | ||
``` | ||
|
||
When rewritten using `bindingCatch` it becomes: | ||
|
||
```kotlin | ||
import arrow.effects.* | ||
import arrow.typeclasses.* | ||
|
||
ForFluxK extensions { | ||
bindingCatch { | ||
val songUrl = getSongUrlAsync().bind() | ||
val musicPlayer = MediaPlayer.load(songUrl) | ||
val totalTime = musicPlayer.getTotaltime() | ||
|
||
val end = DirectProcessor.create<Unit>() | ||
Flux.interval(Duration.ofMillis(100)).takeUntilOther(end).bind() | ||
|
||
val tick = musicPlayer.getCurrentTime().bind() | ||
val percent = (tick / totalTime * 100).toInt() | ||
if (percent >= 100) { | ||
end.onNext(Unit) | ||
} | ||
|
||
percent | ||
}.fix() | ||
} | ||
``` | ||
|
||
Note that any unexpected exception, like `AritmeticException` when `totalTime` is 0, is automatically caught and wrapped inside the flux. | ||
|
||
### Subscription and cancellation | ||
|
||
Flux streams created with comprehensions like `bindingCatch` behave the same way regular flux streams do, including cancellation by disposing the subscription. | ||
|
||
```kotlin | ||
val disposable = | ||
songFlux.value() | ||
.subscribe({ println("Song $it") }, { System.err.println("Error $it") }) | ||
|
||
disposable.dispose() | ||
``` | ||
Note that [`MonadDefer`]({{ '/docs/effects/monaddefer' | relative_url }}) provides an alternative to `bindingCatch` called `bindingCancellable` returning a `arrow.Disposable`. | ||
Invoking this `Disposable` causes an `BindingCancellationException` in the chain which needs to be handled by the subscriber, similarly to what `Deferred` does. | ||
|
||
```kotlin | ||
val (flux, disposable) = | ||
FluxK.monadDefer().bindingCancellable { | ||
val userProfile = Flux.create { getUserProfile("123") } | ||
val friendProfiles = userProfile.friends().map { friend -> | ||
bindDefer { getProfile(friend.id) } | ||
} | ||
listOf(userProfile) + friendProfiles | ||
} | ||
|
||
flux.value() | ||
.subscribe({ println("User $it") }, { System.err.println("Boom! caused by $it") }) | ||
|
||
disposable() | ||
// Boom! caused by BindingCancellationException | ||
``` | ||
|
||
## Available Instances | ||
|
||
* [Applicative]({{ '/docs/typeclasses/applicative' | relative_url }}) | ||
* [ApplicativeError]({{ '/docs/typeclasses/applicativeerror' | relative_url }}) | ||
* [Functor]({{ '/docs/typeclasses/functor' | relative_url }}) | ||
* [Monad]({{ '/docs/typeclasses/monad' | relative_url }}) | ||
* [MonadError]({{ '/docs/typeclasses/monaderror' | relative_url }}) | ||
* [MonadDefer]({{ '/docs/effects/monaddefer' | relative_url }}) | ||
* [Async]({{ '/docs/effects/async' | relative_url }}) | ||
* [Effect]({{ '/docs/effects/effect' | relative_url }}) | ||
* [Foldable]({{ '/docs/typeclasses/foldable' | relative_url }}) | ||
* [Traverse]({{ '/docs/typeclasses/traverse' | relative_url }}) |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,19 @@ | ||
dependencies { | ||
compile project(':arrow-data') | ||
compile project(':arrow-effects') | ||
compile "org.jetbrains.kotlin:kotlin-stdlib-jdk7:$kotlinVersion" | ||
compile project(':arrow-annotations') | ||
kapt project(':arrow-annotations-processor') | ||
kaptTest project(':arrow-annotations-processor') | ||
compileOnly project(':arrow-annotations-processor') | ||
testCompileOnly project(':arrow-annotations-processor') | ||
testCompile "io.kotlintest:kotlintest:$kotlinTestVersion" | ||
testCompile project(':arrow-test') | ||
|
||
compile "io.projectreactor:reactor-core:3.1.3.RELEASE" | ||
testCompile "io.projectreactor:reactor-test:3.1.3.RELEASE" | ||
} | ||
|
||
apply from: rootProject.file('gradle/gradle-mvn-push.gradle') | ||
apply from: rootProject.file('gradle/generated-kotlin-sources.gradle') | ||
apply plugin: 'kotlin-kapt' |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,4 @@ | ||
# Maven publishing configuration | ||
POM_NAME=Arrow-Effects-Reactor | ||
POM_ARTIFACT_ID=arrow-effects-reactor | ||
POM_PACKAGING=jar |
120 changes: 120 additions & 0 deletions
120
modules/effects/arrow-effects-reactor/src/main/kotlin/arrow/effects/FluxK.kt
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,120 @@ | ||
package arrow.effects | ||
|
||
import arrow.Kind | ||
import arrow.core.Either | ||
import arrow.core.Eval | ||
import arrow.core.Left | ||
import arrow.core.Right | ||
import arrow.core.identity | ||
import arrow.effects.typeclasses.Proc | ||
import arrow.higherkind | ||
import arrow.typeclasses.Applicative | ||
import reactor.core.publisher.Flux | ||
import reactor.core.publisher.FluxSink | ||
|
||
fun <A> Flux<A>.k(): FluxK<A> = FluxK(this) | ||
|
||
fun <A> FluxKOf<A>.value(): Flux<A> = | ||
this.fix().flux | ||
|
||
@higherkind | ||
data class FluxK<A>(val flux: Flux<A>) : FluxKOf<A>, FluxKKindedJ<A> { | ||
fun <B> map(f: (A) -> B): FluxK<B> = | ||
flux.map(f).k() | ||
|
||
fun <B> ap(fa: FluxKOf<(A) -> B>): FluxK<B> = | ||
flatMap { a -> fa.fix().map { ff -> ff(a) } } | ||
|
||
fun <B> flatMap(f: (A) -> FluxKOf<B>): FluxK<B> = | ||
flux.flatMap { f(it).fix().flux }.k() | ||
|
||
fun <B> concatMap(f: (A) -> FluxKOf<B>): FluxK<B> = | ||
flux.concatMap { f(it).fix().flux }.k() | ||
|
||
fun <B> switchMap(f: (A) -> FluxKOf<B>): FluxK<B> = | ||
flux.switchMap { f(it).fix().flux }.k() | ||
|
||
fun <B> foldLeft(b: B, f: (B, A) -> B): B = flux.reduce(b, f).block() | ||
|
||
fun <B> foldRight(lb: Eval<B>, f: (A, Eval<B>) -> Eval<B>): Eval<B> { | ||
fun loop(fa_p: FluxK<A>): Eval<B> = when { | ||
fa_p.flux.hasElements().map { !it }.block() -> lb | ||
else -> f(fa_p.flux.blockFirst(), Eval.defer { loop(fa_p.flux.skip(1).k()) }) | ||
} | ||
|
||
return Eval.defer { loop(this) } | ||
} | ||
|
||
fun <G, B> traverse(GA: Applicative<G>, f: (A) -> Kind<G, B>): Kind<G, FluxK<B>> = GA.run { | ||
foldRight(Eval.always { GA.just(Flux.empty<B>().k()) }) { a, eval -> | ||
f(a).map2Eval(eval) { Flux.concat(Flux.just<B>(it.a), it.b.flux).k() } | ||
}.value() | ||
} | ||
|
||
fun handleErrorWith(function: (Throwable) -> FluxK<A>): FluxK<A> = | ||
this.fix().flux.onErrorResume { t: Throwable -> function(t).flux }.k() | ||
|
||
fun runAsync(cb: (Either<Throwable, A>) -> FluxKOf<Unit>): FluxK<Unit> = | ||
flux.flatMap { cb(Right(it)).value() }.onErrorResume { cb(Left(it)).value() }.k() | ||
|
||
companion object { | ||
fun <A> just(a: A): FluxK<A> = | ||
Flux.just(a).k() | ||
|
||
fun <A> raiseError(t: Throwable): FluxK<A> = | ||
Flux.error<A>(t).k() | ||
|
||
operator fun <A> invoke(fa: () -> A): FluxK<A> = | ||
defer { just(fa()) } | ||
|
||
fun <A> defer(fa: () -> FluxKOf<A>): FluxK<A> = | ||
Flux.defer { fa().value() }.k() | ||
|
||
fun <A> runAsync(fa: Proc<A>): FluxK<A> = | ||
Flux.create { emitter: FluxSink<A> -> | ||
fa { either: Either<Throwable, A> -> | ||
either.fold({ | ||
emitter.error(it) | ||
}, { | ||
emitter.next(it) | ||
emitter.complete() | ||
}) | ||
} | ||
}.k() | ||
|
||
tailrec fun <A, B> tailRecM(a: A, f: (A) -> FluxKOf<Either<A, B>>): FluxK<B> { | ||
val either = f(a).fix().value().blockFirst() | ||
return when (either) { | ||
is Either.Left -> tailRecM(either.a, f) | ||
is Either.Right -> Flux.just(either.b).k() | ||
} | ||
} | ||
|
||
fun monadFlat(): FluxKMonadInstance = monad() | ||
|
||
fun monadConcat(): FluxKMonadInstance = object : FluxKMonadInstance { | ||
override fun <A, B> Kind<ForFluxK, A>.flatMap(f: (A) -> Kind<ForFluxK, B>): FluxK<B> = | ||
fix().concatMap { f(it).fix() } | ||
} | ||
|
||
fun monadSwitch(): FluxKMonadInstance = object : FluxKMonadErrorInstance { | ||
override fun <A, B> Kind<ForFluxK, A>.flatMap(f: (A) -> Kind<ForFluxK, B>): FluxK<B> = | ||
fix().switchMap { f(it).fix() } | ||
} | ||
|
||
fun monadErrorFlat(): FluxKMonadErrorInstance = monadError() | ||
|
||
fun monadErrorConcat(): FluxKMonadErrorInstance = object : FluxKMonadErrorInstance { | ||
override fun <A, B> Kind<ForFluxK, A>.flatMap(f: (A) -> Kind<ForFluxK, B>): FluxK<B> = | ||
fix().concatMap { f(it).fix() } | ||
} | ||
|
||
fun monadErrorSwitch(): FluxKMonadErrorInstance = object : FluxKMonadErrorInstance { | ||
override fun <A, B> Kind<ForFluxK, A>.flatMap(f: (A) -> Kind<ForFluxK, B>): FluxK<B> = | ||
fix().switchMap { f(it).fix() } | ||
} | ||
} | ||
} | ||
|
||
inline fun <A, G> FluxKOf<Kind<G, A>>.sequence(GA: Applicative<G>): Kind<G, FluxK<A>> = | ||
fix().traverse(GA, ::identity) |
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Make sure to add this line to the README in the root of the repo. It isn't the same as this one, sadly.