U vi+*@s^ddlmZmZmZmZmZddlmZmZmZddl m Z edZ deej e eej eejeee gee fdddZeje eee gee fd d d Zee ej e ee d d dZeeejdddZeeejdddZeeejdddZeeejdddZejeee gee fdddZd dddddddgZdS))AnyCallableListOptionalTypeVar) Observableabctyping)CompositeDisposable_TN)on_nexton_error on_completedreturncs$ttttdfdd }|S)Nsourcercs4dtjtttjtjdfdd }t|S)a Invokes an action for each element in the observable sequence and invokes an action on graceful or exceptional termination of the observable sequence. This method can be used for debugging, logging, etc. of query behavior by intercepting the message stream to run arbitrary actions for messages on the pipeline. Examples: >>> do_action(send)(observable) >>> do_action(on_next, on_error)(observable) >>> do_action(on_next, on_error, on_completed)(observable) Args: on_next: [Optional] Action to invoke for each element in the observable sequence. on_error: [Optional] Action to invoke on exceptional termination of the observable sequence. on_completed: [Optional] Action to invoke on graceful termination of the observable sequence. Returns: An observable source sequence with the side-effecting behavior applied. Nobserver schedulerrcsRtddfdd }tddfdd }ddfdd }j||||d S) N)xrc sXs|nDz |Wn,tk rH}z|W5d}~XYnX|dSNr Exceptionr )re)rr H/opt/alt/python38/lib/python3.8/site-packages/reactivex/operators/_do.py_on_next,s  zBdo_action_..do_action..subscribe.._on_next exceptionrc sXs|nDz |Wn,tk rH}z|W5d}~XYnX|dSrr r)rr)rr rr _on_error7s  zCdo_action_..do_action..subscribe.._on_errorrc sRsn@z Wn,tk rD}z|W5d}~XYnXdSrrrr )r)rrrr _on_completedBs  zGdo_action_..do_action..subscribe.._on_completedr)r r subscribe)rrrr r#)rr r rrrr%(s   z0do_action_..do_action..subscribe)Nr ObserverBaser r SchedulerBaseDisposableBaserrr%rr r rr do_actions)zdo_action_..do_action)rr )r r rr.rr,r do_action_ s Er/)rrcCst|j|j|jS)a.Invokes an action for each element in the observable sequence and invokes an action on graceful or exceptional termination of the observable sequence. This method can be used for debugging, logging, etc. of query behavior by intercepting the message stream to run arbitrary actions for messages on the pipeline. >>> do(observer) Args: observer: Observer Returns: An operator function that takes the source observable and returns the source sequence with the side-effecting behavior applied. )r/r r rr&rrrdo_Vsr0)r after_nextrcs0dtjtttjtjdfdd }t|S)zInvokes an action with each element after it has been emitted downstream. This can be helpful for debugging, logging, and other side effects. after_next -- Action to invoke on each element after it has been emitted Nrcs&tdfdd }|jjS)N)valuec sHz||Wn,tk rB}z|W5d}~XYnXdSrr)r2r)r1rrrr ws   z1do_after_next..subscribe..on_next)r r%r r)rrr r1rr&rr%tsz do_after_next..subscribe)Nr')rr1r%rr3r do_after_nextks  r4)r on_subscribecs0dtjtttjtjdfdd }t|S)zInvokes an action on subscription. This can be helpful for debugging, logging, and other side effects on the start of an operation. Args: on_subscribe: Action to invoke on subscription Nrcsj|j|j|j|dSNr$)r%r r r)rrr5rrrr%sz"do_on_subscribe..subscribe)Nrr(rrr)r*r)rr5r%rr7rdo_on_subscribes  r9)r on_disposecsFGfdddtjdtjtttjtjdfdd }t|S)zInvokes an action on disposal. This can be helpful for debugging, logging, and other side effects on the disposal of an operation. Args: on_dispose: Action to invoke on disposal cseZdZddfdd ZdS)z do_on_dispose..OnDisposeNr!cs dSrrselfr:rrdisposesz(do_on_dispose..OnDispose.dispose)__name__ __module__ __qualname__r>rr=rr OnDisposesrBNrcs8t}|j|j|j|j|d}|||Sr6)r addr%r r r)rrcomposite_disposable subscription)rBrrrr%s  z do_on_dispose..subscribe)N)rr*r(rrr)r)rr:r%r)rBr:rr do_on_disposes rF)r on_terminatecs0dtjtttjtjdfdd }t|S)a Invokes an action on an on_complete() or on_error() event. This can be helpful for debugging, logging, and other side effects when completion or an error terminates an operation. on_terminate -- Action to invoke when on_complete or throw is called Nrcs6fdd}tdfdd }jj|||dS)Nc sDz Wn,tk r6}z|W5d}~XYn XdSr)rr rerrrrGrrrs  z8do_on_terminate..subscribe..on_completedrc sFz Wn,tk r6}z|W5d}~XYn X|dSr)rr rrIrJrrr s  z4do_on_terminate..subscribe..on_errorr$rr%r rrrr rGrr&rr%sz"do_on_terminate..subscribe)Nr8)rrGr%rrOrdo_on_terminates rP)rafter_terminatecs0dtjtttjtjdfdd }t|S)aInvokes an action after an on_complete() or on_error() event. This can be helpful for debugging, logging, and other side effects when completion or an error terminates an operation on_terminate -- Action to invoke after on_complete or throw is called Nrcs8fdd}tddfdd }jj|||dS)Nc sDz Wn,tk r>}z|W5d}~XYnXdSrr"rHrQrrrrs  z;do_after_terminate..subscribe..on_completedrc sF|z Wn,tk r@}z|W5d}~XYnXdSrrrLrRrrr s   z7do_after_terminate..subscribe..on_errorr$rMrNrQrr&rr%sz%do_after_terminate..subscribe)Nr8)rrQr%rrSrdo_after_terminates rT)finally_actionrcs8Gfdddtjttttdfdd }|S)aInvokes an action after an on_complete(), on_error(), or disposal event occurs. This can be helpful for debugging, logging, and other side effects when completion, an error, or disposal terminates an operation. Note this operator will strive to execute the finally_action once, and prevent any redudant calls Args: finally_action -- Action to invoke after on_complete, on_error, or disposal is called cs0eZdZeedddZddfdd ZdS)zdo_finally..OnDispose was_invokedcSs ||_dSrrV)r<rWrrr__init__sz&do_finally..OnDispose.__init__Nr!cs|jdsd|jd<dSNrTrVr;rUrrr>s z%do_finally..OnDispose.dispose)r?r@rArboolrXr>rrZrrrBsrBrcs2dtjtttjtjdfdd }t|S)Nrcsbdgfdd}tdfdd }t}|jj|||d}|||S)NFc sTzds dd<Wn,tk rN}z|W5d}~XYnXdSrYr"rHrUrrWrrrs zDdo_finally..partial..subscribe..on_completedrKc sV|zds"dd<Wn,tk rP}z|W5d}~XYnXdSrYrrLr\rrr (s  z@do_finally..partial..subscribe..on_errorr$)rr rCr%r )rrrr rDrE)rBrUr)rrWrr%s   z.do_finally..partial..subscribe)Nr'r+rBrUr-rpartials!zdo_finally..partial)rr*rr )rUr^rr]r do_finallys $r_)NNN)r rrrrrZ reactivexrrZreactivex.disposabler r ZOnNextZOnErrorZ OnCompletedr/r(r0r4Actionr9rFrPrTr___all__rrrrsB   M( #" B