Commit cb3f9b97 authored by topjohnwu's avatar topjohnwu

More tweaking to Rx pipeline

parent 97843532
...@@ -7,7 +7,7 @@ import com.topjohnwu.magisk.model.entity.module.Repo ...@@ -7,7 +7,7 @@ import com.topjohnwu.magisk.model.entity.module.Repo
@Dao @Dao
abstract class RepoDao { abstract class RepoDao {
val repoIDSet: Set<String> get() = getRepoID().map { it.id }.toSet() val repoIDList get() = getRepoID().map { it.id }
val repos: List<Repo> get() = when (Config.repoOrder) { val repos: List<Repo> get() = when (Config.repoOrder) {
Config.Value.ORDER_NAME -> getReposNameOrder() Config.Value.ORDER_NAME -> getReposNameOrder()
...@@ -45,7 +45,7 @@ abstract class RepoDao { ...@@ -45,7 +45,7 @@ abstract class RepoDao {
abstract fun removeRepo(id: String) abstract fun removeRepo(id: String)
@Query("DELETE FROM repos WHERE id IN (:idList)") @Query("DELETE FROM repos WHERE id IN (:idList)")
abstract fun removeRepos(idList: List<String>) abstract fun removeRepos(idList: Collection<String>)
@Query("SELECT * FROM etag") @Query("SELECT * FROM etag")
protected abstract fun etagRaw(): RepoEtag? protected abstract fun etagRaw(): RepoEtag?
......
...@@ -7,6 +7,7 @@ import com.topjohnwu.magisk.data.network.GithubApiServices ...@@ -7,6 +7,7 @@ import com.topjohnwu.magisk.data.network.GithubApiServices
import com.topjohnwu.magisk.model.entity.module.Repo import com.topjohnwu.magisk.model.entity.module.Repo
import io.reactivex.Flowable import io.reactivex.Flowable
import io.reactivex.Single import io.reactivex.Single
import io.reactivex.rxkotlin.toFlowable
import io.reactivex.schedulers.Schedulers import io.reactivex.schedulers.Schedulers
import se.ansman.kotshi.JsonSerializable import se.ansman.kotshi.JsonSerializable
import timber.log.Timber import timber.log.Timber
...@@ -21,16 +22,13 @@ class RepoUpdater( ...@@ -21,16 +22,13 @@ class RepoUpdater(
) { ) {
private fun loadRepos(repos: List<GithubRepoInfo>, cached: MutableSet<String>) = private fun loadRepos(repos: List<GithubRepoInfo>, cached: MutableSet<String>) =
Flowable.fromIterable(repos).map { repos.toFlowable().parallel().runOn(Schedulers.io()).map {
it.id to dateFormat.parse(it.pushed_at)!!
}.parallel().runOn(Schedulers.io()).map {
// Skip submission // Skip submission
if (it.first == "submission") if (it.id == "submission")
return@map return@map
(repoDB.getRepo(it.first)?.apply { (repoDB.getRepo(it.id)?.apply { cached.remove(it.id) } ?:
cached.remove(it.first) Repo(it.id)).runCatching {
} ?: Repo(it.first)).runCatching { update(it.pushDate)
update(it.second)
repoDB.addRepo(this) repoDB.addRepo(this)
}.getOrElse { Timber.e(it) } }.getOrElse { Timber.e(it) }
}.sequential() }.sequential()
...@@ -57,8 +55,8 @@ class RepoUpdater( ...@@ -57,8 +55,8 @@ class RepoUpdater(
} }
} }
private fun forcedReload(cached: MutableSet<String>) = Flowable.fromIterable(cached) private fun forcedReload(cached: MutableSet<String>) =
.parallel().runOn(Schedulers.io()).map { cached.toFlowable().parallel().runOn(Schedulers.io()).map {
runCatching { runCatching {
Repo(it).update() Repo(it).update()
}.getOrElse { Timber.e(it) } }.getOrElse { Timber.e(it) }
...@@ -67,9 +65,9 @@ class RepoUpdater( ...@@ -67,9 +65,9 @@ class RepoUpdater(
private fun String.trimEtag() = substring(indexOf('\"'), lastIndexOf('\"') + 1) private fun String.trimEtag() = substring(indexOf('\"'), lastIndexOf('\"') + 1)
operator fun invoke(forced: Boolean = false) : Single<Unit> { operator fun invoke(forced: Boolean = false) : Single<Unit> {
val cached = Collections.synchronizedSet(HashSet(repoDB.repoIDSet)) val cached = Collections.synchronizedSet(HashSet(repoDB.repoIDList))
return loadPage(cached, etag = repoDB.etagKey).doOnComplete { return loadPage(cached, etag = repoDB.etagKey).doOnComplete {
repoDB.removeRepos(cached.toList()) repoDB.removeRepos(cached)
}.onErrorResumeNext { it: Throwable -> }.onErrorResumeNext { it: Throwable ->
if (it is CachedException) { if (it is CachedException) {
if (forced) if (forced)
...@@ -78,21 +76,22 @@ class RepoUpdater( ...@@ -78,21 +76,22 @@ class RepoUpdater(
Timber.e(it) Timber.e(it)
} }
Flowable.empty() Flowable.empty()
}.collect({}, {_, _ -> }) }.ignoreElements().toSingleDefault(Unit)
}
companion object {
private val dateFormat: SimpleDateFormat =
SimpleDateFormat("yyyy-MM-dd'T'HH:mm:ss'Z'", Locale.US).apply {
timeZone = TimeZone.getTimeZone("UTC")
}
} }
object CachedException : Exception() object CachedException : Exception()
} }
private val dateFormat: SimpleDateFormat =
SimpleDateFormat("yyyy-MM-dd'T'HH:mm:ss'Z'", Locale.US).apply {
timeZone = TimeZone.getTimeZone("UTC")
}
@JsonSerializable @JsonSerializable
data class GithubRepoInfo( data class GithubRepoInfo(
@Json(name = "name") val id: String, @Json(name = "name") val id: String,
val pushed_at: String val pushed_at: String
) ) {
\ No newline at end of file @Transient
val pushDate = dateFormat.parse(pushed_at)!!
}
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment