@@ -36,16 +36,9 @@ import pekko.stream.Attributes.{ InputBuffer, LogLevels }
3636import pekko .stream .Attributes .SourceLocation
3737import pekko .stream .OverflowStrategies ._
3838import pekko .stream .Supervision .Decider
39- import pekko .stream .impl .{
40- Buffer => BufferImpl ,
41- ContextPropagation ,
42- FailedSource ,
43- JavaStreamSource ,
44- ReactiveStreamsCompliance ,
45- TraversalBuilder
46- }
39+ import pekko .stream .impl .{ Buffer => BufferImpl , ContextPropagation , ReactiveStreamsCompliance , TraversalBuilder }
4740import pekko .stream .impl .Stages .DefaultAttributes
48- import pekko .stream .impl .fusing .GraphStages .{ FutureSource , SimpleLinearGraphStage , SingleSource }
41+ import pekko .stream .impl .fusing .GraphStages .SimpleLinearGraphStage
4942import pekko .stream .scaladsl .{
5043 DelayStrategy ,
5144 Source ,
@@ -2169,36 +2162,19 @@ private[pekko] object TakeWithin {
21692162 override def onPull (): Unit = pull(in)
21702163
21712164 @ nowarn(" msg=Any" )
2172- @ tailrec
21732165 def onFailure (ex : Throwable ): Unit = {
21742166 import Collect .NotApplied
21752167 if (maximumRetries < 0 || attempt < maximumRetries) {
21762168 pf.applyOrElse(ex, NotApplied ) match {
21772169 case _ : NotApplied .type => failStage(ex)
21782170 case source : Graph [SourceShape [T ] @ unchecked, M @ unchecked] if TraversalBuilder .isEmptySource(source) =>
21792171 completeStage()
2180- case source : Graph [SourceShape [T ] @ unchecked, M @ unchecked] =>
2181- TraversalBuilder .getValuePresentedSource(source) match {
2182- case OptionVal .Some (graph) => graph match {
2183- case singleSource : SingleSource [T @ unchecked] => emit(out, singleSource.elem, () => completeStage())
2184- case failed : FailedSource [T @ unchecked] => onFailure(failed.failure)
2185- case futureSource : FutureSource [T @ unchecked] => futureSource.future.value match {
2186- case Some (Success (elem)) => emit(out, elem, () => completeStage())
2187- case Some (Failure (ex)) => onFailure(ex)
2188- case None =>
2189- switchTo(source)
2190- attempt += 1
2191- }
2192- case iterableSource : IterableSource [T @ unchecked] =>
2193- emitMultiple(out, iterableSource.elements, () => completeStage())
2194- case javaStreamSource : JavaStreamSource [T @ unchecked, _] =>
2195- emitMultiple(out, javaStreamSource.open().spliterator(), () => completeStage())
2196- case _ =>
2197- switchTo(source)
2198- attempt += 1
2199- }
2172+ case other : Graph [SourceShape [T ] @ unchecked, M @ unchecked] =>
2173+ TraversalBuilder .getSingleSource(other) match {
2174+ case OptionVal .Some (singleSource) =>
2175+ emit(out, singleSource.elem.asInstanceOf [T ], () => completeStage())
22002176 case _ =>
2201- switchTo(source )
2177+ switchTo(other )
22022178 attempt += 1
22032179 }
22042180 case _ => throw new IllegalStateException () // won't happen, compiler exhaustiveness check pleaser
0 commit comments