U vi@sddlZddlZddlmZddlmZmZmZddlm Z mZddl m Z m Z m Z ddlmZedZed ZGd d d eZdS) N)Future)ListOptionalTypeVar)abctyping)CompositeDisposable DisposableSingleAssignmentDisposable)AsyncIOScheduler_TStateZRxc@seZdZdZdejeeeej dddZ dej ejeeeej dddZ dej ejeeeej dd d Zed d d ZdS)AsyncIOThreadSafeSchedulerzA scheduler that schedules work via the asyncio mainloop. This is a subclass of AsyncIOScheduler which uses the threadsafe asyncio methods. N)actionstatereturncsLtddfdd }j|ddfdd }tt|S)a!Schedules an action to be executed. Args: action: Action to be executed. state: [Optional] state to be given to the action function. Returns: The disposable object used to cancel the scheduled action (best effort). Nrcsjd_dSNrZ invoke_actionZ disposablersadselfrri/opt/alt/python38/lib/python3.8/site-packages/reactivex/scheduler/eventloop/asynciothreadsafescheduler.pyinterval(sz5AsyncIOThreadSafeScheduler.schedule..intervalcsFrdStddfdd }j|dS)NrcsddSNr)cancel set_resultr)futurehandlerr cancel_handle4szKAsyncIOThreadSafeScheduler.schedule..dispose..cancel_handle)_on_self_loop_or_not_runningrr_loopcall_soon_threadsaferesultr!r r)rrdispose-s z4AsyncIOThreadSafeScheduler.schedule..dispose)r r#r$rr )rrrrr(r)rr rrrrschedules  z#AsyncIOThreadSafeScheduler.schedule)duetimerrrcs|dkr jdStddfdd gddfdd }j|ddfd d }tt|S) auSchedules an action to be executed after duetime. Args: duetime: Relative time after which to execute the action. action: Action to be executed. state: [Optional] state to be given to the action function. Returns: The disposable object used to cancel the scheduled action (best effort). rrNrcsjd_dSrrrrrrrTsz>AsyncIOThreadSafeScheduler.schedule_relative..intervalcsjdSN)appendr# call_laterr)r rsecondsrrrstage2[sz.stage2csVddfdd r$dStddfdd }j|dS)Nrcs6zWntk r0YnXdSr+)popr Exceptionr)r rrdo_cancel_handlesas  zXAsyncIOThreadSafeScheduler.schedule_relative..dispose..do_cancel_handlescsddSr)rrr2rrrr!nszTAsyncIOThreadSafeScheduler.schedule_relative..dispose..cancel_handle)r"rr#r$r%r&r'r3rr(`s z=AsyncIOThreadSafeScheduler.schedule_relative..dispose)Z to_secondsr)r r,r#r$rr )rr*rrr/r(r)rr rrr.rrrschedule_relative=s z,AsyncIOThreadSafeScheduler.schedule_relativecCs ||}|j||j||dS)aoSchedules an action to be executed at duetime. Args: duetime: Absolute time at which to execute the action. action: Action to be executed. state: [Optional] state to be given to the action function. Returns: The disposable object used to cancel the scheduled action (best effort). r) to_datetimer4now)rr*rrrrrschedule_absolutews z,AsyncIOThreadSafeScheduler.schedule_absolutercCs>|jsdSd}z t}Wntk r2YnX|j|kS)z Returns True if either self._loop is not running, or we're currently executing on self._loop. In both cases, waiting for a future to be resolved on the loop would result in a deadlock. TN)r# is_runningasyncioget_event_loop RuntimeError)r current_looprrrr"s  z7AsyncIOThreadSafeScheduler._on_self_loop_or_not_running)N)N)N)__name__ __module__ __qualname____doc__rZScheduledActionr rrZDisposableBaser)Z RelativeTimer4Z AbsoluteTimer7boolr"rrrrrs, ( > r)r9loggingconcurrent.futuresrrrrrZ reactivexrZreactivex.disposablerr r Zasyncioschedulerr r getLoggerlogrrrrrs