Reaktive Documentation

repository·master·Indexed 22 days ago

https://github.com/badoo/reaktive

A Kotlin Multiplatform implementation of Reactive Extensions (Rx) providing a consistent API for reactive programming across JVM, Android, iOS, JavaScript, and Native platforms. Includes modules for annotations, testing, and interoperability with Kotlin Coroutines and RxJava (v2 and v3). Features include built-in schedulers, DisposableScope for subscription lifecycle management, a Plugin API for decorating sources, and specific guidance for Swift interop using wrappers and RxSwift bridging.

Tokens
5.1K
Snippets
12
Records
18
Agent score
79%

What's inside Reaktive

  1. Set up the iOS sample application

    master

    To run the iOS sample application, you must first install the CocoaPods dependencies.

    1. Run pod install in the project directory to install all required Pod dependencies.
    2. Open the Xcode project.
    3. The framework is assembled automatically during the Xcode build process using the embedAndSignAppleFrameworkForXcode Gradle task. For more details on how this task works, refer to the Kotlin Multiplatform documentation.
    pod install
  2. Manage subscriptions with DisposableScope

    master

    DisposableScope provides an easy way to manage the lifecycle of multiple subscriptions. When the scope is disposed, all subscriptions and disposables associated with it are automatically cleaned up.

    You can use the DisposableScope interface via delegation to integrate it directly into your classes (like Presenters or Activities).

    Usage Example

    val scope = disposableScope {
        observable.subscribeScoped(...) // Disposed when scope is disposed
    
        doOnDispose { /* Called when scope is disposed */ }
    
        someDisposable.scope() // someDisposable is disposed when scope is disposed
    }
    
    // Later
    scope.dispose()

    Integration Example

    class MyPresenter(...) : DisposableScope by DisposableScope() {
        fun load() {
            // Subscription will be disposed when the presenter is disposed
            longRunningAction.subscribeScoped(onComplete = view::hideProgressBar)
        }
    }
    
    class MyActivity : AppCompatActivity(), DisposableScope by DisposableScope() {
        override fun onCreate(savedInstanceState: Bundle?) {
            super.onCreate(savedInstanceState)
            MyPresenter(...).scope()
        }
    
        override fun onDestroy() {
            dispose()
            super.onDestroy()
        }
    }
    val scope =
        disposableScope {
            observable.subscribeScoped(...) // Subscription will be disposed when the scope is disposed
    
            doOnDispose {
                // Will be called when the scope is disposed
            }
    
            someDisposable.scope() // `someDisposable` will be disposed when the scope is disposed
        }
    
    // At some point later
    scope.dispose()
  3. Export Reaktive to Swift via Gradle

    master

    To use Reaktive in a Swift project, you must export the Reaktive library within your Kotlin Multiplatform Gradle configuration. This involves adding the Reaktive dependency as an api dependency in commonMain and explicitly exporting it in the framework block for iOS targets.

    kotlin {
        targets
            .filterIsInstance<KotlinNativeTarget>()
            .filter { it.konanTarget.family == Family.IOS }
            .forEach { target ->
                target.binaries {
                    framework {
                        // Some setup code here
    
                        export("com.badoo.reaktive:reaktive:<version>")
                    }
                }
            }
    
        sourceSets {
            named("commonMain") {
                dependencies {
                    api("com.badoo.reaktive:reaktive:<version>")
                }
            }
        }
    }
  4. Interop between Reaktive and RxSwift

    master

    While there is no official interop module, you can bridge Reaktive wrappers to RxSwift by implementing extension methods on RxSwift.Observable, RxSwift.Single, RxSwift.Maybe, and RxSwift.Completable. These extensions use the .create operator to map Reaktive subscription events to RxSwift observer events.

    import RxSwift
    import YourKotlinFramework
    
    extension RxSwift.Observable where Element : AnyObject {
        static func from(_ observable: ObservableWrapper<Element>) -> RxSwift.Observable<Element> {
            return RxSwift.Observable<Element>.create { observer in
                let disposable = observable.subscribe(
                    onError: { observer.onError(KotlinError($0)) },
                    onComplete: observer.onCompleted,
                    onNext: observer.onNext
                )
    
                return Disposables.create(with: disposable.dispose)
            }
        }
    }
    
    extension RxSwift.Single where Element : AnyObject {
        static func from(_ single: SingleWrapper<Element>) -> RxSwift.Single<Element> {
            return RxSwift.Single<Element>.create { observer in
                let disposable = single.subscribe(
                    onError: { observer(.failure(KotlinError($0))) },
                    onSuccess: { observer(.success($0)) }
                )
    
                return Disposables.create(with: disposable.dispose)
            }
        }
    }
    
    extension RxSwift.Maybe where Element : AnyObject {
        static func from(_ maybe: MaybeWrapper<Element>) -> RxSwift.Maybe<Element> {
            return RxSwift.Maybe<Element>.create { observer in
                let disposable = maybe.subscribe(
                    onError: { observer(.error(KotlinError($0))) },
                    onComplete: { observer(.completed) },
                    onSuccess: { observer(.success($0)) }
                )
    
                return Disposables.create(with: disposable.dispose)
            }
        }
    }
    
    extension RxSwift.Completable {
        static func from(_ completable: CompletableWrapper) -> RxSwift.Completable {
            return RxSwift.Completable.create { observer in
                let disposable = completable.subscribe(
                    onError: { observer(.error(KotlinError($0))) },
                    onComplete: { observer(.completed) }
                )
    
                return Disposables.create(with: disposable.dispose)
            }
        }
    }
    
    struct KotlinError : Error {
        let throwable: KotlinThrowable
    
        init (_ throwable: KotlinThrowable) {
            self.throwable = throwable
        }
    }
  5. Install Reaktive and its modules

    master

    Reaktive is distributed via Maven Central as several multiplatform modules. Depending on your needs, you can include the core library, annotations, testing utilities, or interoperability modules for Coroutines and RxJava.

    Add the following to your build.gradle file, replacing <version> with the desired version number:

    kotlin {
        sourceSets {
            commonMain {
                dependencies {
                    implementation 'com.badoo.reaktive:reaktive:<version>'
                    implementation 'com.badoo.reaktive:reaktive-annotations:<version>'
                    implementation 'com.badoo.reaktive:coroutines-interop:<version>' // For interop with coroutines
                    implementation 'com.badoo.reaktive:rxjava2-interop:<version>' // For interop with RxJava v2
                    implementation 'com.badoo.reaktive:rxjava3-interop:<version>' // For interop with RxJava v3
                }
            }
    
            commonTest {
                dependencies {
                    implementation 'com.badoo.reaktive:reaktive-testing:<version>'
                }
            }
        }
    }
  6. How to use Coroutines interop

    master

    The coroutines-interop module allows you to bridge Reaktive and Kotlin Coroutines. You can convert between Flow and Observable, convert suspend functions to Reaktive types, and bridge CoroutineContext with Reaktive Schedulers.

    Conversion Examples

    // Flow <-> Observable
    val flow: Flow<Int> = observableOf(1, 2, 3).asFlow()
    val observable: Observable<Int> = flowOf(1, 2, 3).asObservable()
    
    // Suspend function -> Single
    fun doSomething() {
        singleFromCoroutine { getSomething() }
            .subscribe { println(it) }
    }
    
    suspend fun getSomething(): String {
        delay(1.seconds)
        return "something"
    }
    
    // Scheduler <-> Dispatcher
    val defaultScheduler = Dispatchers.Default.asScheduler()
    val computationDispatcher = computationScheduler.asCoroutineDispatcher()
    val flow: Flow<Int> = observableOf(1, 2, 3).asFlow()
    val observable: Observable<Int> = flowOf(1, 2, 3).asObservable()
  7. Use Reaktive wrappers in Swift

    master

    Once exposed via wrappers, you can interact with Reaktive sources in Swift using standard subscription patterns. Remember to manage the lifecycle by calling .dispose() on the returned disposable.

    func foo() {
        let ds = SharedDataSource()
    
        let disposable = ds.load().subscribe { result in
            // Handle the result
        }
    
        // At some point later
        disposable.dispose()
    }
  8. Implement a Reaktive Plugin

    master

    The Plugin API allows you to decorate Reaktive sources (Observables, Singles, Maybes, Completables). To create a plugin, implement the ReaktivePlugin interface and register it using registerReaktivePlugin.

    object MyPlugin : ReaktivePlugin {
        override fun <T> onAssembleObservable(observable: Observable<T>): Observable<T> = 
            // Return a decorated Observable
    
        override fun <T> onAssembleSingle(single: Single<T>): Single<T> = TODO()
        override fun <T> onAssembleMaybe(maybe: Maybe<T>): Maybe<T> = TODO()
        override fun onAssembleCompletable(completable: Completable): Completable = TODO()
    }
    object MyPlugin : ReaktivePlugin {
        override fun <T> onAssembleObservable(observable: Observable<T>): Observable<T> =
            object : Observable<T> {
                private val traceException = TraceException()
    
                override fun subscribe(observer: ObservableObserver<T>) {
                    observable.subscribe(
                        object : ObservableObserver<T> by observer {
                            override fun onError(error: Throwable) {
                                observer.onError(error, traceException)
                            }
                        }
                    )
                }
            }
    
        override fun <T> onAssembleSingle(single: Single<T>): Single<T> =
            TODO("Similar to onAssembleSingle")
    
        override fun <T> onAssembleMaybe(maybe: Maybe<T>): Maybe<T> =
            TODO("Similar to onAssembleSingle")
    
        override fun onAssembleCompletable(completable: Completable): Completable =
            TODO("Similar to onAssembleSingle")
    
        private fun ErrorCallback.onError(error: Throwable, traceException: TraceException) {
            if (error.suppressedExceptions.lastOrNull() !is TraceException) {
                error.addSuppressed(traceException)
            }
            onError(error)
        }
    
        private class TraceException : Exception()
    }
  9. Expose Reaktive sources to Swift using wrappers

    master

    Because Kotlin interface generics are not exported to Swift, you cannot use Observable<T>, Single<T>, etc., directly in Swift. Instead, you must use the provided wrap() extension functions in your Kotlin code to convert these sources into wrapper classes that Swift can consume.

    Available Wrappers:

    • Observable<T>.wrap() $\rightarrow$ ObservableWrapper<T>
    • BehaviorObservable<T>.wrap() $\rightarrow$ BehaviorObservableWrapper<T> (also works for BehaviorSubject)
    • Single<T>.wrap() $\rightarrow$ SingleWrapper<T>
    • Maybe<T>.wrap() $\rightarrow$ MaybeWrapper<T>
    • Completable.wrap() $\rightarrow$ CompletableWrapper
    class SharedDataSource {
        fun load(): SingleWrapper<String> =
            singleFromFunction {
                // A long running operation
                "A result"
            }
                .subscribeOn(ioScheduler)
                .observeOn(mainScheduler)
                .wrap()
    }
    
    class SharedViewModel {
        private val _state = BehaviorSubject(State())
        val state: BehaviorObservableWrapper<State> = _state.wrap()
    }
  10. Available Reaktive Schedulers

    master

    Reaktive provides several built-in schedulers to manage execution contexts:

    • computationScheduler: Fixed thread pool equal to the number of cores.
    • ioScheduler: Unbound thread pool with a caching policy.
    • newThreadScheduler: Creates a new thread for each unit of work.
    • singleScheduler: Executes tasks on a single shared background thread.
    • trampolineScheduler: Queues tasks and executes them on one of the participating threads.
    • mainScheduler: Executes tasks on the main thread.