[collab/github/gitlab] chore: cleanup, rename and clarify functions

GitOrigin-RevId: fa9cac621b225fbd1a8eb602828484e8c542a919
This commit is contained in:
Ivan Semenov
2025-09-11 14:45:50 +00:00
committed by intellij-monorepo-bot
parent 07f526c9cb
commit 0a56a2c205
13 changed files with 92 additions and 144 deletions
@@ -4,7 +4,6 @@
package com.intellij.collaboration.async
import com.intellij.collaboration.util.ComputedResult
import com.intellij.collaboration.util.HashingUtil
import com.intellij.openapi.Disposable
import com.intellij.openapi.diagnostic.Logger
import com.intellij.openapi.extensions.ExtensionPointListener
@@ -13,7 +12,6 @@ import com.intellij.openapi.extensions.PluginDescriptor
import com.intellij.openapi.util.Disposer
import com.intellij.platform.util.coroutines.childScope
import com.intellij.util.cancelOnDispose
import com.intellij.util.containers.HashingStrategy
import com.intellij.util.containers.toArray
import com.intellij.util.diff.Diff
import kotlinx.coroutines.*
@@ -304,81 +302,6 @@ fun <T> Flow<T>.stateInNow(cs: CoroutineScope, defaultValue: T): StateFlow<T> {
fun <T> Flow<T>.modelFlow(cs: CoroutineScope, log: Logger): SharedFlow<T> =
catch { log.error(it) }.shareIn(cs, SharingStarted.Lazily, 1)
/**
* The destructor is never necessary because cleanup can be performed on scope cancellation
* @see associateCachingBy
*/
@ApiStatus.Obsolete
fun <T, K, V> Flow<Iterable<T>>.associateCachingBy(
keyExtractor: (T) -> K,
hashingStrategy: HashingStrategy<K>,
valueExtractor: CoroutineScope.(T) -> V,
destroy: suspend V.() -> Unit,
update: (suspend V.(T) -> Unit)? = null,
)
: Flow<Map<K, V>> = flow {
coroutineScope {
val container = MappingScopedItemsContainer(this, keyExtractor, hashingStrategy, valueExtractor, destroy, update)
collect {
container.update(it)
emit(container.mappingState.value)
}
awaitCancellation()
}
}
/**
* Associate each *item* [T] *key* [K] in the iterable from the receiver flow (source list) with a *value* [V]
*
* Keys are distinguished by a [hashingStrategy]
*
* When a new iterable is received:
* * a new [CoroutineScope] and a new value is created via [valueExtractor] for new items
* * existing values are updated via [update] if it was supplied
* * values for missing items are removed and their scope is cancelled
*
* Order of the values in the resulting map is the same as in the source iterable
* All [CoroutineScope]'s of values are only active while the resulting flow is being collected
*
* **Returned flow never completes**
*/
fun <T, K, V> Flow<Iterable<T>>.associateCachingBy(
keyExtractor: (T) -> K,
hashingStrategy: HashingStrategy<K>,
valueExtractor: CoroutineScope.(T) -> V,
update: (suspend V.(T) -> Unit)? = null,
)
: Flow<Map<K, V>> = associateCachingBy(keyExtractor, hashingStrategy, valueExtractor, { }, update)
/**
* @see associateCachingBy
*
* Shorthand for cases where key is the same as item destructor simply cancels the value scope
*/
private fun <T, R> Flow<Iterable<T>>.associateCaching(
hashingStrategy: HashingStrategy<T>,
mapper: CoroutineScope.(T) -> R,
update: (suspend R.(T) -> Unit)? = null,
): Flow<Map<T, R>> {
return associateCachingBy({ it }, hashingStrategy, { mapper(it) }, { }, update)
}
/**
* Creates a list of model objects from DTOs
*/
fun <T, R> Flow<Iterable<T>>.mapDataToModel(
sourceIdentifier: (T) -> Any,
mapper: CoroutineScope.(T) -> R,
update: (suspend R.(T) -> Unit),
): Flow<List<R>> =
associateCaching(HashingUtil.mappingStrategy(sourceIdentifier), mapper, update).map { it.values.toList() }
/**
* Create a list of view models from models
*/
fun <T, R> Flow<Iterable<T>>.mapModelsToViewModels(mapper: CoroutineScope.(T) -> R): Flow<List<R>> =
associateCaching(HashingStrategy.identity(), mapper).map { it.values.toList() }
/**
* Maps each item in the collection from the source flow to a flow and emits the array of the latest values of each mapped flow.
* Each new emission of the source flow triggers re-subscription to the mapped flows.
@@ -498,29 +421,6 @@ fun <T, R> Flow<ComputedResult<T>>.transformConsecutiveSuccesses(
}
}
/**
* Transforms the flow of some computation requests to a flow of computation states of this request
* Will not emit "loading" state if the computation was completed before handling its state
*/
@OptIn(ExperimentalCoroutinesApi::class)
@ApiStatus.Internal
fun <T> Flow<Deferred<T>>.computationState(): Flow<ComputedResult<T>> =
transformLatest { request ->
if (!request.isCompleted) {
emit(ComputedResult.loading())
}
try {
val value = request.await()
emit(ComputedResult.success(value))
}
catch (e: Exception) {
if (e !is CancellationException) {
emit(ComputedResult.failure(e))
}
}
}
@OptIn(ExperimentalCoroutinesApi::class)
inline fun <A, T> computationStateFlow(arguments: Flow<A>, crossinline computer: suspend (A) -> T): Flow<ComputedResult<T>> =
arguments.transformLatest { parameters ->
@@ -1,19 +1,75 @@
// Copyright 2000-2024 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license.
package com.intellij.collaboration.async
import com.intellij.collaboration.util.HashingUtil
import com.intellij.platform.util.coroutines.childScope
import com.intellij.util.containers.CollectionFactory
import com.intellij.util.containers.HashingStrategy
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.NonCancellable
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.updateAndGet
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withContext
import org.jetbrains.annotations.ApiStatus
/**
* Associate each *item* [T] *key* [K] in the iterable from the receiver flow (source list) with a *value* [V]
*
* Keys are distinguished by a [hashingStrategy]
*
* When a new iterable is received:
* * a new [CoroutineScope] and a new value is created via [valueExtractor] for new items
* * existing values are updated via [update] if it was supplied
* * values for missing items are removed and their scope is cancelled
*
* Order of the values in the resulting map is the same as in the source iterable
* All [CoroutineScope]'s of values are only active while the resulting flow is being collected
*
* **Returned flow never completes**
*/
fun <T, K, V> Flow<Iterable<T>>.associateCachingBy(
keyExtractor: (T) -> K,
hashingStrategy: HashingStrategy<K>,
valueExtractor: CoroutineScope.(T) -> V,
update: (suspend V.(T) -> Unit)? = null,
): Flow<Map<K, V>> = flow {
coroutineScope {
val container = MappingScopedItemsContainer(this, keyExtractor, hashingStrategy, valueExtractor, update)
collect {
container.update(it)
emit(container.mappingState.value)
}
awaitCancellation()
}
}
/**
* @see associateCachingBy
*/
fun <T, R> Flow<Iterable<T>>.associateCachingWith(
hashingStrategy: HashingStrategy<T>,
mapper: CoroutineScope.(T) -> R,
update: (suspend R.(T) -> Unit)? = null,
): Flow<Map<T, R>> {
return associateCachingBy({ it }, hashingStrategy, { mapper(it) }, update)
}
/**
* Creates a list of stateful objects from the list of DTO (Data Transfer Object)
* Stateful objects are updated with [update] when a DTO identified by a [sourceIdentifier] changes
*/
fun <T, R> Flow<Iterable<T>>.mapDataToModel(
sourceIdentifier: (T) -> Any,
mapper: CoroutineScope.(T) -> R,
update: (suspend R.(T) -> Unit),
): Flow<List<R>> =
associateCachingWith(HashingUtil.mappingStrategy(sourceIdentifier), mapper, update).map { it.values.toList() }
/**
* Creates a list of stateful objects from other stateful objects, comparing the original objects by identity
*/
fun <T, R> Flow<Iterable<T>>.mapStatefulToStateful(mapper: CoroutineScope.(T) -> R): Flow<List<R>> =
associateCachingWith(HashingStrategy.identity(), mapper).map { it.values.toList() }
/**
* Allows mapping a collection of items [T] to scoped (coroutine scope bound) values [V]
* An intermittent key [K] is used to uniquely identify items
@@ -31,7 +87,6 @@ class MappingScopedItemsContainer<T, K, V> internal constructor(
private val keyExtractor: (T) -> K,
private val hashingStrategy: HashingStrategy<K>,
private val mapper: CoroutineScope.(T) -> V,
private val destroy: suspend V.() -> Unit,
private val update: (suspend V.(T) -> Unit)? = null,
) {
private val _mappingState = MutableStateFlow<Map<K, ScopingWrapper<V>>>(emptyMap())
@@ -63,7 +118,6 @@ class MappingScopedItemsContainer<T, K, V> internal constructor(
val deletedKeys = currentMap.keys - resultMap.keys
for (key in deletedKeys) {
val scopedValue = currentMap[key] ?: continue
scopedValue.value.destroy()
scopedValue.cancel()
}
@@ -85,10 +139,10 @@ class MappingScopedItemsContainer<T, K, V> internal constructor(
}
companion object {
fun <T, V> byIdentity(cs: CoroutineScope, mapper: CoroutineScope.(T) -> V) =
fun <T, V> byIdentity(cs: CoroutineScope, mapper: CoroutineScope.(T) -> V): MappingScopedItemsContainer<T, T?, V> =
MappingScopedItemsContainer(cs, { it }, HashingStrategy.identity(), mapper, {})
fun <T, V> byEquality(cs: CoroutineScope, mapper: CoroutineScope.(T) -> V) =
fun <T, V> byEquality(cs: CoroutineScope, mapper: CoroutineScope.(T) -> V): MappingScopedItemsContainer<T, T?, V> =
MappingScopedItemsContainer(cs, { it }, HashingStrategy.canonical(), mapper, {})
}
}
@@ -152,7 +152,7 @@ internal class GHPRDiffViewModelImpl(
}.map { it.getOrNull().orEmpty() }.stateInNow(cs, emptyMap())
private val mappedThreads: StateFlow<List<MappedGHPRReviewThreadDiffViewModel>> =
threadsVm.compactThreads.mapModelsToViewModels { sharedVm ->
threadsVm.compactThreads.mapStatefulToStateful { sharedVm ->
MappedGHPRReviewThreadDiffViewModel(this, sharedVm, threadMappings.mapNotNull { it[sharedVm.id] })
}.stateInNow(cs, emptyList())
@@ -27,7 +27,6 @@ import com.intellij.openapi.util.Key
import com.intellij.platform.util.coroutines.childScope
import com.intellij.util.cancelOnDispose
import com.intellij.util.concurrency.annotations.RequiresEdt
import com.intellij.util.containers.HashingStrategy
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.flow.MutableStateFlow
@@ -104,11 +103,9 @@ internal class GHPRReviewDiffExtension : DiffExtension() {
GHPRReviewUnifiedPosition(change, leftLine, rightLine)
}.apply {
cs.launchNow {
inlays.associateCachingBy(
keyExtractor = { it },
hashingStrategy = HashingStrategy.identity(),
valueExtractor = { inlay -> GHPRInlayUtils.installInlayHoverOutline(this, editor, side, locationToLine, inlay) }
).collect()
inlays
.mapStatefulToStateful { inlayModel -> GHPRInlayUtils.installInlayHoverOutline(this, editor, side, locationToLine, inlayModel) }
.collect()
}
}
}
@@ -132,9 +129,9 @@ private class DiffEditorModel(
@RequiresEdt private val lineToUnified: (Int) -> GHPRReviewUnifiedPosition,
) : GHPRReviewDiffEditorModel {
private val threads = diffVm.threads.mapModelsToViewModels { MappedThread(cs, it) }.stateInNow(cs, emptyList())
private val newComments = diffVm.newComments.mapModelsToViewModels { MappedNewComment(it) }.stateInNow(cs, emptyList())
private val aiComments = diffVm.aiComments.mapModelsToViewModels { MappedAIComment(it) }.stateInNow(cs, emptyList())
private val threads = diffVm.threads.mapStatefulToStateful { MappedThread(cs, it) }.stateInNow(cs, emptyList())
private val newComments = diffVm.newComments.mapStatefulToStateful { MappedNewComment(it) }.stateInNow(cs, emptyList())
private val aiComments = diffVm.aiComments.mapStatefulToStateful { MappedAIComment(it) }.stateInNow(cs, emptyList())
override val inlays: StateFlow<Collection<GHPREditorMappedComponentModel>> =
combineStateIn(cs, threads, newComments, aiComments) { threads, new, ai -> threads + new + ai }
@@ -2,8 +2,8 @@
package org.jetbrains.plugins.github.pullrequest.ui.editor
import com.intellij.collaboration.async.combineState
import com.intellij.collaboration.async.mapModelsToViewModels
import com.intellij.collaboration.async.mapState
import com.intellij.collaboration.async.mapStatefulToStateful
import com.intellij.collaboration.async.stateInNow
import com.intellij.collaboration.ui.codereview.editor.*
import com.intellij.collaboration.util.ExcludingApproximateChangedRangesShifter
@@ -58,8 +58,8 @@ internal class GHPRReviewFileEditorModel internal constructor(
}.stateInNow(cs, null)
override val inlays: StateFlow<Collection<GHPREditorMappedComponentModel>> = combine(
fileVm.threads.mapModelsToViewModels { ShiftedThread(it) },
fileVm.newComments.mapModelsToViewModels { ShiftedNewComment(it) },
fileVm.threads.mapStatefulToStateful { ShiftedThread(it) },
fileVm.newComments.mapStatefulToStateful { ShiftedNewComment(it) },
) { threads, new ->
// very explicit ordering: if we order back to front, loading of editor appears smoother (most initial loading happens off-screen)
threads.sortedByDescending { it.line.value ?: -1 } + new
@@ -126,7 +126,7 @@ internal class GHPRReviewFileEditorViewModelImpl(
allMappedThreads.mapState { map -> map.filterValues { it.change == change } }
override val threads: StateFlow<Collection<GHPRReviewFileEditorThreadViewModel>> =
threadsVm.compactThreads.mapModelsToViewModels { sharedVm ->
threadsVm.compactThreads.mapStatefulToStateful { sharedVm ->
MappedGHPRReviewEditorThreadViewModel(this, sharedVm, mappedThreads.mapNotNull { it[sharedVm.id] })
}.stateInNow(cs, emptyList())
@@ -1,10 +1,10 @@
// Copyright 2000-2024 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license.
package org.jetbrains.plugins.github.pullrequest.ui.editor
import com.intellij.collaboration.async.associateCachingBy
import com.intellij.collaboration.async.collectScoped
import com.intellij.collaboration.async.launchNow
import com.intellij.collaboration.async.mapScoped
import com.intellij.collaboration.async.mapStatefulToStateful
import com.intellij.collaboration.ui.codereview.diff.DiscussionsViewOption
import com.intellij.collaboration.ui.codereview.editor.*
import com.intellij.collaboration.util.HashingUtil
@@ -25,7 +25,6 @@ import com.intellij.openapi.fileEditor.OpenFileDescriptor
import com.intellij.openapi.project.Project
import com.intellij.openapi.util.Disposer
import com.intellij.util.cancelOnDispose
import com.intellij.util.containers.HashingStrategy
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import org.jetbrains.plugins.github.pullrequest.config.GithubPullRequestsProjectUISettings
@@ -123,11 +122,9 @@ private suspend fun showReview(project: Project, settings: GithubPullRequestsPro
val userIcon = fileVm.iconProvider.getIcon(fileVm.currentUser.url, 16)
editor.renderInlays(model.inlays, HashingUtil.mappingStrategy(GHPREditorMappedComponentModel::key)) {
launchNow {
model.inlays.associateCachingBy(
keyExtractor = { it },
hashingStrategy = HashingStrategy.identity(),
valueExtractor = { inlay -> GHPRInlayUtils.installInlayHoverOutline(this, editor, Side.RIGHT, null, inlay) }
).collect()
model.inlays
.mapStatefulToStateful { inlayModel -> GHPRInlayUtils.installInlayHoverOutline(this, editor, Side.RIGHT, null, inlayModel) }
.collect()
}
createRenderer(it, userIcon)
}
@@ -99,9 +99,9 @@ private class DiffEditorModel(
) : CodeReviewEditorModel<GitLabMergeRequestEditorMappedComponentModel> {
override val inlays: StateFlow<Collection<GitLabMergeRequestEditorMappedComponentModel>> = combine(
diffVm.discussions.mapModelsToViewModels { MappedDiscussion(it) },
diffVm.draftDiscussions.mapModelsToViewModels { MappedDraftNote(it) },
diffVm.newDiscussions.mapModelsToViewModels { MappedNewDiscussion(it) }
diffVm.discussions.mapStatefulToStateful { MappedDiscussion(it) },
diffVm.draftDiscussions.mapStatefulToStateful { MappedDraftNote(it) },
diffVm.newDiscussions.mapStatefulToStateful { MappedNewDiscussion(it) }
) { discussions, drafts, new ->
discussions + drafts + new
}.stateInNow(cs, emptyList())
@@ -3,7 +3,7 @@ package org.jetbrains.plugins.gitlab.mergerequest.ui.editor
import com.intellij.collaboration.async.combineState
import com.intellij.collaboration.async.launchNow
import com.intellij.collaboration.async.mapModelsToViewModels
import com.intellij.collaboration.async.mapStatefulToStateful
import com.intellij.collaboration.async.stateInNow
import com.intellij.collaboration.ui.codereview.editor.*
import com.intellij.collaboration.util.ExcludingApproximateChangedRangesShifter
@@ -63,9 +63,9 @@ internal class GitLabMergeRequestEditorReviewUIModel internal constructor(
}.stateInNow(cs, null)
override val inlays: StateFlow<Collection<GitLabMergeRequestEditorMappedComponentModel>> = combine(
fileVm.discussions.mapModelsToViewModels { ShiftedDiscussion(it) },
fileVm.draftNotes.mapModelsToViewModels { ShiftedDraftNote(it) },
fileVm.newDiscussions.mapModelsToViewModels { ShiftedNewDiscussion(it) }
fileVm.discussions.mapStatefulToStateful { ShiftedDiscussion(it) },
fileVm.draftNotes.mapStatefulToStateful { ShiftedDraftNote(it) },
fileVm.newDiscussions.mapStatefulToStateful { ShiftedNewDiscussion(it) }
) { discussions, drafts, new ->
discussions + drafts + new
}.stateInNow(cs, emptyList())
@@ -62,13 +62,13 @@ internal class GitLabMergeRequestDiscussionsViewModelsImpl(
override val discussions: DiscussionsFlow = mergeRequest.discussions
.throwFailure()
.mapModelsToViewModels { GitLabMergeRequestDiscussionViewModelBase(project, this, projectData, currentUser, it) }
.mapStatefulToStateful { GitLabMergeRequestDiscussionViewModelBase(project, this, projectData, currentUser, it) }
.modelFlow(cs, LOG)
override val draftNotes: DraftNotesFlow = mergeRequest.draftNotes
.throwFailure()
.mapFiltered { it.discussionId == null }
.mapModelsToViewModels { GitLabMergeRequestStandaloneDraftNoteViewModelBase(project, this, it, mergeRequest) }
.mapStatefulToStateful { GitLabMergeRequestStandaloneDraftNoteViewModelBase(project, this, it, mergeRequest) }
.modelFlow(cs, LOG)
@@ -78,7 +78,7 @@ class GitLabMergeRequestTimelineDiscussionViewModelImpl(
override val replies: StateFlow<List<GitLabNoteViewModel>> = discussion.notes
.map { it.drop(1) }
.mapModelsToViewModels { GitLabNoteViewModelImpl(project, this, projectData, it, flowOf(false), currentUser) }
.mapStatefulToStateful { GitLabNoteViewModelImpl(project, this, projectData, it, flowOf(false), currentUser) }
.stateIn(cs, SharingStarted.Lazily, listOf())
override val isBusy: StateFlow<Boolean> = taskLauncher.busy
@@ -3,7 +3,7 @@ package org.jetbrains.plugins.gitlab.ui.clone.model
import com.intellij.collaboration.api.HttpStatusErrorException
import com.intellij.collaboration.async.collectBatches
import com.intellij.collaboration.async.mapModelsToViewModels
import com.intellij.collaboration.async.mapStatefulToStateful
import com.intellij.collaboration.async.withInitial
import com.intellij.collaboration.messages.CollaborationToolsBundle
import com.intellij.openapi.components.service
@@ -119,7 +119,7 @@ internal class GitLabCloneRepositoriesListViewModelImpl(
@OptIn(ExperimentalCoroutinesApi::class)
private val listsPerAccount = reloadSignal.withInitial(Unit).flatMapLatest { _ ->
accountManager.accountsState.mapModelsToViewModels<GitLabAccount, GitLabCloneRepositoriesForAccountViewModel> { account ->
accountManager.accountsState.mapStatefulToStateful<GitLabAccount, GitLabCloneRepositoriesForAccountViewModel> { account ->
GitLabCloneRepositoriesForAccountViewModelImpl(this, accountManager, account)
}
}.stateIn(cs, SharingStarted.Eagerly, listOf())
@@ -61,7 +61,7 @@ internal class GitLabMergeRequestDiscussionViewModelBase(
}.stateInNow(cs, null)
private val initialNotesSize: Int = discussion.notes.value.size
private val notesVms = discussion.notes.mapModelsToViewModels { note ->
private val notesVms = discussion.notes.mapStatefulToStateful { note ->
GitLabNoteViewModelImpl(project, this, projectData, note, discussion.notes.map { it.firstOrNull()?.id == note.id }, currentUser)
}.stateInNow(cs, emptyList())
override val notes: StateFlow<List<NoteItem>> =