File tree 1 file changed +14
-1
lines changed
reactive/kotlinx-coroutines-reactive/src
1 file changed +14
-1
lines changed Original file line number Diff line number Diff line change @@ -6,6 +6,7 @@ package kotlinx.coroutines.reactive
6
6
import kotlinx.atomicfu.*
7
7
import kotlinx.coroutines.*
8
8
import kotlinx.coroutines.channels.*
9
+ import kotlinx.coroutines.intrinsics.*
9
10
import kotlinx.coroutines.selects.*
10
11
import kotlinx.coroutines.sync.*
11
12
import org.reactivestreams.*
@@ -104,10 +105,22 @@ public class PublisherCoroutine<in T>(
104
105
// registerSelectSend
105
106
@Suppress(" PARAMETER_NAME_CHANGED_ON_OVERRIDE" )
106
107
override fun <R > registerSelectClause2 (select : SelectInstance <R >, element : T , block : suspend (SendChannel <T >) -> R ) {
107
- mutex.onLock.registerSelectClause2(select, null ) {
108
+ val clause = suspend {
108
109
doLockedNext(element)?.let { throw it }
109
110
block(this )
110
111
}
112
+
113
+ // TODO discuss it
114
+ launch(start = CoroutineStart .UNDISPATCHED ) {
115
+ mutex.lock()
116
+ // Already selected -- bail out
117
+ if (! select.trySelect()) {
118
+ mutex.unlock()
119
+ return @launch
120
+ }
121
+
122
+ clause.startCoroutineCancellable(select.completion)
123
+ }
111
124
}
112
125
113
126
/*
You can’t perform that action at this time.
0 commit comments