This commit is contained in:
Horis
2024-01-31 23:38:38 +08:00
parent 6dd9f83f35
commit 0ef03dba9a
7 changed files with 121 additions and 51 deletions
@@ -2,6 +2,7 @@ package io.legado.app.constant
import androidx.annotation.IntDef import androidx.annotation.IntDef
@Suppress("ConstPropertyName")
object BookSourceType { object BookSourceType {
const val default = 0 // 0 文本 const val default = 0 // 0 文本
@@ -250,8 +250,12 @@ interface BookSourceDao {
@get:Query("select * from book_sources where loginUrl is not null and loginUrl != ''") @get:Query("select * from book_sources where loginUrl is not null and loginUrl != ''")
val allLogin: List<BookSource> val allLogin: List<BookSource>
@get:Query("select * from book_sources where enabled = 1 and bookSourceType = 0 order by customOrder") @get:Query(
val allTextEnabled: List<BookSource> """select bookSourceUrl, bookSourceName, bookSourceGroup, customOrder, enabled, enabledExplore,
trim(loginUrl) <> '' hasLoginUrl, lastUpdateTime, respondTime, weight, trim(exploreUrl) <> '' hasExploreUrl
from book_sources where enabled = 1 and bookSourceType = 0 order by customOrder"""
)
val allTextEnabledPart: List<BookSourcePart>
@get:Query("select distinct bookSourceGroup from book_sources where trim(bookSourceGroup) <> ''") @get:Query("select distinct bookSourceGroup from book_sources where trim(bookSourceGroup) <> ''")
val allGroupsUnProcessed: List<String> val allGroupsUnProcessed: List<String>
@@ -479,6 +479,15 @@ object BookHelp {
} }
} }
fun getDurChapter(
oldBook: Book,
newChapterList: List<BookChapter>
): Int {
return oldBook.run {
getDurChapter(durChapterIndex, durChapterTitle, newChapterList, totalChapterNum)
}
}
private val chapterNamePattern1 by lazy { private val chapterNamePattern1 by lazy {
Pattern.compile(".*?第([\\d零〇一二两三四五六七八九十百千万壹贰叁肆伍陆柒捌玖拾佰仟]+)[章节篇回集话]") Pattern.compile(".*?第([\\d零〇一二两三四五六七八九十百千万壹贰叁肆伍陆柒捌玖拾佰仟]+)[章节篇回集话]")
} }
@@ -10,14 +10,13 @@ import io.legado.app.exception.NoStackTraceException
import io.legado.app.help.config.AppConfig import io.legado.app.help.config.AppConfig
import io.legado.app.ui.book.search.SearchScope import io.legado.app.ui.book.search.SearchScope
import io.legado.app.utils.getPrefBoolean import io.legado.app.utils.getPrefBoolean
import io.legado.app.utils.mapNotNullParallel import io.legado.app.utils.mapParallelSafe
import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.ExecutorCoroutineDispatcher import kotlinx.coroutines.ExecutorCoroutineDispatcher
import kotlinx.coroutines.Job import kotlinx.coroutines.Job
import kotlinx.coroutines.asCoroutineDispatcher import kotlinx.coroutines.asCoroutineDispatcher
import kotlinx.coroutines.currentCoroutineContext import kotlinx.coroutines.currentCoroutineContext
import kotlinx.coroutines.ensureActive import kotlinx.coroutines.ensureActive
import kotlinx.coroutines.flow.buffer
import kotlinx.coroutines.flow.catch import kotlinx.coroutines.flow.catch
import kotlinx.coroutines.flow.collect import kotlinx.coroutines.flow.collect
import kotlinx.coroutines.flow.flow import kotlinx.coroutines.flow.flow
@@ -83,16 +82,11 @@ class SearchModel(private val scope: CoroutineScope, private val callBack: CallB
} }
}.onStart { }.onStart {
callBack.onSearchStart() callBack.onSearchStart()
}.mapNotNullParallel(threadCount) { }.mapParallelSafe(threadCount) {
try { withTimeout(30000L) {
withTimeout(30000L) { WebBook.searchBookAwait(it, searchKey, searchPage)
WebBook.searchBookAwait(it, searchKey, searchPage)
}
} catch (e: Throwable) {
currentCoroutineContext().ensureActive()
null
} }
}.buffer(0).onEach { items -> }.onEach { items ->
appDb.searchBookDao.insert(*items.toTypedArray()) appDb.searchBookDao.insert(*items.toTypedArray())
mergeItems(items, precision) mergeItems(items, precision)
currentCoroutineContext().ensureActive() currentCoroutineContext().ensureActive()
@@ -114,7 +114,7 @@ class CheckSourceService : BaseService() {
upNotification() upNotification()
}.onEachParallel(threadCount) { }.onEachParallel(threadCount) {
checkSource(it) checkSource(it)
}.buffer(0).onEach { }.onEach {
finishCount++ finishCount++
notificationMsg = getString( notificationMsg = getString(
R.string.progress_show, R.string.progress_show,
@@ -34,10 +34,19 @@ import io.legado.app.ui.book.searchContent.SearchResult
import io.legado.app.utils.DocumentUtils import io.legado.app.utils.DocumentUtils
import io.legado.app.utils.FileUtils import io.legado.app.utils.FileUtils
import io.legado.app.utils.isContentScheme import io.legado.app.utils.isContentScheme
import io.legado.app.utils.mapParallelSafe
import io.legado.app.utils.postEvent import io.legado.app.utils.postEvent
import io.legado.app.utils.toStringArray import io.legado.app.utils.toStringArray
import io.legado.app.utils.toastOnUi import io.legado.app.utils.toastOnUi
import kotlinx.coroutines.Dispatchers.IO import kotlinx.coroutines.Dispatchers.IO
import kotlinx.coroutines.flow.catch
import kotlinx.coroutines.flow.collect
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.flow.onCompletion
import kotlinx.coroutines.flow.onEach
import kotlinx.coroutines.flow.onEmpty
import kotlinx.coroutines.flow.onStart
import kotlinx.coroutines.flow.take
import java.io.File import java.io.File
import java.io.FileInputStream import java.io.FileInputStream
import java.io.FileNotFoundException import java.io.FileNotFoundException
@@ -266,39 +275,44 @@ class ReadBookViewModel(application: Application) : BaseViewModel(application) {
private fun autoChangeSource(name: String, author: String) { private fun autoChangeSource(name: String, author: String) {
if (!AppConfig.autoChangeSource) return if (!AppConfig.autoChangeSource) return
execute { execute {
val sources = appDb.bookSourceDao.allTextEnabled val sources = appDb.bookSourceDao.allTextEnabledPart
sources.forEach { source -> flow {
WebBook.preciseSearchAwait(this, source, name, author).getOrNull()?.let { book -> for (source in sources) {
if (book.tocUrl.isEmpty()) { source.getBookSource()?.let {
WebBook.getBookInfoAwait(source, book) emit(it)
}
val toc = WebBook.getChapterListAwait(source, book).getOrThrow()
val chapter = toc.getOrElse(book.durChapterIndex) {
toc.last()
}
val nextChapter = toc.getOrElse(chapter.index) {
toc.first()
}
kotlin.runCatching {
WebBook.getContentAwait(
bookSource = source,
book = book,
bookChapter = chapter,
nextChapterUrl = nextChapter.url
)
changeTo(book, toc)
return@execute
} }
} }
} }.onStart {
throw NoStackTraceException("没有合适书源") ReadBook.upMsg(context.getString(R.string.source_auto_changing))
}.onStart { }.mapParallelSafe(AppConfig.threadCount) { source ->
ReadBook.upMsg(context.getString(R.string.source_auto_changing)) val book = WebBook.preciseSearchAwait(this, source, name, author).getOrThrow()
}.onError { if (book.tocUrl.isEmpty()) {
AppLog.put("自动换源失败\n${it.localizedMessage}", it) WebBook.getBookInfoAwait(source, book)
context.toastOnUi("自动换源失败\n${it.localizedMessage}") }
}.onFinally { val toc = WebBook.getChapterListAwait(source, book).getOrThrow()
ReadBook.upMsg(null) val chapter = toc.getOrElse(book.durChapterIndex) {
toc.last()
}
val nextChapter = toc.getOrElse(chapter.index) {
toc.first()
}
WebBook.getContentAwait(
bookSource = source,
book = book,
bookChapter = chapter,
nextChapterUrl = nextChapter.url
)
book to toc
}.take(1).onEach { (book, toc) ->
changeTo(book, toc)
}.onEmpty {
throw NoStackTraceException("没有合适书源")
}.onCompletion {
ReadBook.upMsg(null)
}.catch {
AppLog.put("自动换源失败\n${it.localizedMessage}", it)
context.toastOnUi("自动换源失败\n${it.localizedMessage}")
}.collect()
} }
} }
@@ -2,7 +2,11 @@ package io.legado.app.utils
import kotlinx.coroutines.ExperimentalCoroutinesApi import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.async import kotlinx.coroutines.async
import kotlinx.coroutines.currentCoroutineContext
import kotlinx.coroutines.ensureActive
import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.FlowCollector
import kotlinx.coroutines.flow.buffer
import kotlinx.coroutines.flow.channelFlow import kotlinx.coroutines.flow.channelFlow
import kotlinx.coroutines.flow.filterNotNull import kotlinx.coroutines.flow.filterNotNull
import kotlinx.coroutines.flow.flatMapMerge import kotlinx.coroutines.flow.flatMapMerge
@@ -20,13 +24,57 @@ inline fun <T> Flow<T>.onEachParallel(
action(value) action(value)
emit(value) emit(value)
} }
} }.buffer(0)
@OptIn(ExperimentalCoroutinesApi::class)
inline fun <T> Flow<T>.onEachParallelSafe(
concurrency: Int,
crossinline action: suspend (T) -> Unit
): Flow<T> = flatMapMerge(concurrency) { value ->
flow {
try {
action(value)
} catch (e: Throwable) {
currentCoroutineContext().ensureActive()
}
emit(value)
}
}.buffer(0)
@OptIn(ExperimentalCoroutinesApi::class) @OptIn(ExperimentalCoroutinesApi::class)
inline fun <T, R> Flow<T>.mapParallel( inline fun <T, R> Flow<T>.mapParallel(
concurrency: Int, concurrency: Int,
crossinline transform: suspend (T) -> R, crossinline transform: suspend (T) -> R,
): Flow<R> = flatMapMerge(concurrency) { value -> flow { emit(transform(value)) } } ): Flow<R> = flatMapMerge(concurrency) { value -> flow { emit(transform(value)) } }.buffer(0)
@OptIn(ExperimentalCoroutinesApi::class)
inline fun <T, R> Flow<T>.mapParallelSafe(
concurrency: Int,
crossinline transform: suspend (T) -> R,
): Flow<R> = flatMapMerge(concurrency) { value ->
flow {
try {
emit(transform(value))
} catch (e: Throwable) {
currentCoroutineContext().ensureActive()
}
}
}.buffer(0)
@OptIn(ExperimentalCoroutinesApi::class)
inline fun <T, R> Flow<T>.transformParallelSafe(
concurrency: Int,
crossinline transform: suspend FlowCollector<R>.(T) -> R,
): Flow<R> = flatMapMerge(concurrency) { value ->
flow {
try {
transform(value)
} catch (e: Throwable) {
currentCoroutineContext().ensureActive()
}
}
}.buffer(0)
inline fun <T, R> Flow<T>.mapNotNullParallel( inline fun <T, R> Flow<T>.mapNotNullParallel(
concurrency: Int, concurrency: Int,
@@ -67,7 +115,7 @@ inline fun <T, R> Flow<T>.mapAsync(
}.map { }.map {
it.await() it.await()
}.onEach { semaphore.release() } }.onEach { semaphore.release() }
} }.buffer(0)
} }
inline fun <T, R> Flow<T>.mapAsyncIndexed( inline fun <T, R> Flow<T>.mapAsyncIndexed(
@@ -89,7 +137,7 @@ inline fun <T, R> Flow<T>.mapAsyncIndexed(
}.map { }.map {
it.await() it.await()
}.onEach { semaphore.release() } }.onEach { semaphore.release() }
} }.buffer(0)
} }
inline fun <T> Flow<T>.onEachAsync( inline fun <T> Flow<T>.onEachAsync(
@@ -110,7 +158,7 @@ inline fun <T> Flow<T>.onEachAsync(
}.map { }.map {
it.await() it.await()
}.onEach { semaphore.release() } }.onEach { semaphore.release() }
} }.buffer(0)
} }
inline fun <T> Flow<T>.onEachAsyncIndexed( inline fun <T> Flow<T>.onEachAsyncIndexed(
@@ -135,5 +183,5 @@ inline fun <T> Flow<T>.onEachAsyncIndexed(
}.map { }.map {
it.await() it.await()
}.onEach { semaphore.release() } }.onEach { semaphore.release() }
} }.buffer(0)
} }