U vi@sddlmZddlmZmZmZmZddlZddlmZm Z ddl m Z m Z m Z ddlmZedZeeedeeegeeffeed d d Zd gZdS) )Future)CallableOptionalTypeVarUnionN) Observableabc)CompositeDisposableSerialDisposableSingleAssignmentDisposable)CurrentThreadScheduler_Tz Future[_T])sourcesreturncs6t|dtjtttjtjdfdd }t|S)aQContinues an observable sequence that is terminated normally or by an exception with the next observable sequence. Examples: >>> res = reactivex.on_error_resume_next(xs, ys, zs) Returns: An observable sequence that concatenates the source sequences, even if a sequence terminates exceptionally. N)observer schedulerrcsR|p t}tt}dtjttddfdd ||_t |S)N)rstatercsz t}Wntk r*YdSXt|r<||n|}t|trTt|n|}t}|_ dt t ddfdd }|j j ||d|_ dS)N)rrcs|dS)N)schedule)r)actionrW/opt/alt/python38/lib/python3.8/site-packages/reactivex/observable/onerrorresumenext.py on_resume<szKon_error_resume_next_..subscribe..action..on_resumer)N)next StopIterationZ on_completedcallable isinstancer reactivexZ from_futurer disposabler Exception subscribeZon_next)rrsourcecurrentdr)rrsources_ subscriptionrrr*s" z8on_error_resume_next_..subscribe..action)N) r Z singletonr r SchedulerBaserrrrr )rrZ cancelabler$)rrr%rr "s  z(on_error_resume_next_..subscribe)N)iterrZ ObserverBaser rr&ZDisposableBaser)rr rr'ron_error_resume_next_s$r))asynciortypingrrrrrrrZreactivex.disposabler r r Zreactivex.schedulerr r rr)__all__rrrrs   9