How to make the function return Observable

Viewed 288

I am trying to make my code more Rx, so here is problem that i run into during refactor:

I had previously function like this:

private fun startLocationUpdates() {
    compositeDisposable.add(rxLocationRepository.onLocationUpdate()
            .subscribe(Consumer { location ->
                appendGeoEvent(location)
            }))

    rxScheduleStopLocationUpdates()
}

Which i turn this into something like this:

val locationUpdates = rxLocationRepository.onLocationUpdate()
        .concatMap(appendGeoEvent()) 
        .doOnSubscribe(rxScheduleStopLocationUpdates()) 
        .publish()


private fun startLocationUpdates() {compositeDisposable.add(locationUpdates.connect()) } 

The problem is now with the appendGeoEvent to make work with my created locationUpdates ConnectableObservable.

The appendGeoEvent looked like this(it was two connected functions)

 fun appendGeoEvent(location: Location) {
    val locationGeoEvent = LocationGeoEvent(
            accuracy = location.accuracy.toDouble(),
            latitude = location.latitude,
            longitude = location.longitude,
            timestampGeoEvent = location.time
    )
    appendGeoEvent(locationGeoEvent)
}

private fun appendGeoEvent(event: LocationGeoEvent) {
    synchronized(mEventsMonitor) {
        if (!normalSaveScheduled.get()) {
            configureNormalSaveCompletable()
        }
        locationEvents.add(event)
    }
}

So now i need to make one function of this two that will return some Observable or Copmletable. Again, how to make it have sense and work with this locationUpdates ConnectableObservable:

val locationUpdates = rxLocationRepository.onLocationUpdate()
    .concatMap(appendGeoEvent()) 
    .doOnSubscribe(rxScheduleStopLocationUpdates()) 
    .publish()

The other part of code is this:

1)The configureNormalSaveCompletable method:

 fun configureNormalSaveCompletable() {
    scheduleNormalSave.subscribe(Action { rxSaveEvents(locationEvents, false) })
}

2) The scheduleNormalSave Completable(it schedules action of doing something in future and sets the boolean needed in rxSaveEvents method)

    val scheduleNormalSave = Completable.timer(locationEventsSaveThrottlingSeconds, TimeUnit.SECONDS).doOnSubscribe(Consumer { normalSaveScheduled.set(true) }).doOnComplete(Action { normalSaveScheduled.set(false) })

3) And here is the rxSaveEvents method:

    private fun rxSaveEvents(
        events: List<LocationEvent>,
        locationUpdatesAboutToStop: Boolean): Boolean {

    synchronized(mEventsMonitor) {
        if (events.isNotEmpty()) {
            val locationEvents = LocationEvents()
            segregateLocationEvents(events, locationEvents)
            scoreReadings(locationEvents)
            sendLocationEventsInternal(locationEvents, locationUpdatesAboutToStop)

            Completable.fromCallable { updateLocationEventsToDb(locationEvents) }.subscribeOn(Schedulers.io()).subscribe()

        }
        if (locationUpdatesAboutToStop) {
            stopLocationUpdates()

            Completable.fromCallable { cleanAllReadingsData() }.subscribeOn(Schedulers.io()).subscribe()
            Completable.fromCallable { removeExpiredData() }.subscribeOn(Schedulers.io()).subscribe()
        }

        locationEvents.clear()

        return true
    }
}

UPDATE

What i've tried so far is this:

//what should be in the method parameter and how to exctract this 
    fun apendGeoEvent(location: Location): Observable<LocationGeoEvent>? {
    return Observable.just(location)
            .flatMap { t: Location ->  Observable.just(LocationGeoEvent(
                accuracy = location.accuracy.toDouble(),
                    latitude= location.latitude,
                    longitude = location.longitude,
                    timestampGeoEvent = location.time
                )
            )
            }
            .doOnNext{
                if (!normalSaveScheduled.get()) {
                    configureNormalSaveCompletable()
                }
                locationEvents.add(it)
            }
}

However i don't know how to take that Observable from the upstream to my appendGeoEvent method:

    val locationUpdates = rxLocationRepository.onLocationUpdate()
        .concatMap(appendGeoEvent()) 
        .doOnSubscribe(rxScheduleStopLocationUpdates()) 
        .publish()
0 Answers
Related