a Kf= @sxdZddlZddlZddlZddlZddlZddlZddlZddlZddl m Z m Z m Z m Z mZmZmZmZmZmZddlZddlmZddlmZddlmZddlmZddlmZdd lmZdd lmZdd lmZdd lmZdd lm Z ddlm!Z!ddlm"Z"ddlm#Z#ddl$Ze%e&Z'd(ej)Z*dZ+e,dduZ-ej.j/ej.j0ej.j1ej.j2ej.j3ej.j4fZ5ej.j/ej.j0ej.j1ej.j2ej.j4fZ6ej.j/ej.j2ej.j3ej.j4fZ7ej.j/ej.j2ej.j4fZ8dZ9dZ:dZ;ee<ee<dddZ=eej>ee?e?dddZ@GdddeAZBeBej>e?dddd ZCejDeBeeee d!d"d#ZEeBeee#d$d%d&ZFe eBeejGejHfe"ee#dd'd(d)ZIe?eBe?d*d+d,ZJGd-d.d.ejKejLejMZNGd/d0d0ejKejOZPGd1d2d2ePejLejMZQGd3d4d4ePejLejMZRe ee<e"eee<eeSeejKfd5d6d7ZTeBejHeUee<ee!ee!ejLffd8d9d:ZVeeeWeeejXd;dd?ZZee<ee<d@dAdBZ[GdCdDdDej\Z]GdEdFdFej^Z_GdGdHdHej^Z`GdIdJdJejaZbGdKdLdLejcZdGdMdNdNeWZeGdOdPdPeAZfefddQdRdSZgefdTdUdVZhGdWdXdXeAZieiee ejjgdfdQdYdZZkeiejjee ejjgdfdd[d\d]Zleiee ejjgdfdd^d_d`ZmeiejneUddadbdcZoeie ejjgdfeUddddedfZpeie ejjgdfddgdhdiZqeeeejreedjdkdlZseeeeeeefdmdndoZtGdpdqdqejnZndS)rz.Invocation-side implementation of gRPC Python.N) AnyCallableDictIteratorListOptionalSequenceSetTupleUnion)_common) _compression)_grpcio_metadata)_observability)cygrpc)ChannelArgumentType)DeserializingFunction)IntegratedCallFactory) MetadataType)NullaryCallbackType) ResponseType)SerializingFunction)UserTagzgrpc-python/{}!GRPC_SINGLE_THREADED_UNARY_STREAMz0Exception calling channel subscription callback!z?<{} of RPC that terminated with: status = {} details = "{}" >zZ<{} of RPC that terminated with: status = {} details = "{}" debug_error_string = "{}" >timeoutreturncCs|dur dSt|SN)timerr q/sparta/input/_build_configuration/image_build+validate/lib/bmcenv/lib64/python3.9/site-packages/grpc/_channel.py _deadlinemsr")unknown_cygrpc_codedetailsrcCs d||S)Nz,Server sent unknown code {} and details "{}")format)r#r$r r r!_unknown_code_detailsqsr&c@seZdZUejed<eejed<e e ed<e ed<e e ed<e e j ed<e eed<e eed<eed <eeed <e eed <e eed <e eed <e eed<e eed<eeje e e e e e j e edddZddZdS) _RPCState conditiondueinitial_metadataresponsetrailing_metadatacoder$debug_error_string cancelled callbacks fork_epochrpc_start_time rpc_end_timemethodtarget)r)r*r,r-r$cCsjt|_t||_||_d|_||_||_||_ d|_ d|_ d|_ d|_ d|_d|_g|_t|_dSNF) threading Conditionr(setr)r*r+r,r-r$r.r2r3r4r5r/r0rget_fork_epochr1)selfr)r*r,r-r$r r r!__init__s  z_RPCState.__init__cCst|_dSr)r7r8r(r;r r r!reset_postfork_childsz_RPCState.reset_postfork_childN)__name__ __module__ __qualname__r7r8__annotations__r r OperationTyperrrgrpc StatusCodestrboolrrintfloatrr<r>r r r r!r'ys,             (r')stater-r$rcCs0|jdur,||_||_|jdur&d|_d|_dSNr )r-r$r*r,)rJr-r$r r r!_aborts   rL)eventrJresponse_deserializerrc Cs$g}|jD]}|}|j||tjjkr<||_q |tjjkr| }|durt ||}|durd}t |t jj|n||_q |tjjkr ||_|jdurt j|} | durt jj|_t| ||_n| |_||_||_t|_t|| |j!d|_!q |S)Nz!Exception deserializing response!)"batch_operationstyper)removerrCreceive_initial_metadatar*receive_messagemessager deserializerLrDrEINTERNALr+receive_status_on_clientr,r-!CYGRPC_STATUS_CODE_TO_STATUS_CODEgetUNKNOWNr&r$ error_stringr.r perf_counterr3rmaybe_record_rpc_latencyextendr0) rMrJrNr0batch_operationoperation_typeserialized_responser+r$r-r r r! _handle_eventsF              rb)rJrNrcsfdd}|S)Nc sj.t|}jj }Wdn1s:0Y|D]L}z |WqHty}z$tdt|jt|WYd}~qHd}~00qH|oj t kS)NzException in callback %s: %s) r(rb notify_allr) Exceptionloggingerrorreprfuncr1rr:)rMr0donecallbackerNrJr r! handle_events  & z$_event_handler..handle_eventr )rJrNrmr rlr!_event_handlersrn)request_iteratorrJcallrequest_serializer event_handlerrcs6fdd}tj|d}|d|dS)z'Consume a request supplied by the user.csnd}zztt}Wnty@YW|s8tqYnbtytd}tjj}d}t | t j ||t||YW|stdS0W|stn|st0t |}jjdurƈjs|dur.tjj}d} t j ||t||WddSjtjjt|tf}|}|s~jtjjWddSfdd}t jjj|ttjdjdurWddSnWddSWdq1s0YqjZjdurJjtjj t!tf}|}|sJjtjj Wdn1s`0YdS)NFTzException iterating requests!Exception serializing request!csjduptjjjvSr)r-rrC send_messager)r rJr r!_done=s  zJ_consume_request_iterator..consume_request_iterator.._done)spin_cb)"renter_user_request_generatornext StopIteration"return_from_user_request_generatorrdrDrErZ_LOGGER exceptioncancelr !STATUS_CODE_TO_CYGRPC_STATUS_CODErL serializer(r-r/rVr)addrCrtSendMessageOperation _EMPTY_FLAGSoperaterQwait functoolspartialblock_if_fork_in_progresssend_close_from_clientSendCloseFromClientOperation)*return_from_user_request_generator_invokedrequestr-r$serialized_request operations operatingrvrprrrorqrJr r!consume_request_iterator s                2  z;_consume_request_iterator..consume_request_iteratorr5TNrForkManagedThread setDaemonstart)rorJrprqrrrconsumption_threadr rr!_consume_request_iterators O r) class_name rpc_statercCs|j|jdur*d|WdS|jtjjurXt||j|jWdSt||j|j|j WdSWdn1s0YdS)z Calculates error string for RPC.Nz <{} object>) r(r-r%rDrEOK_OK_RENDEZVOUS_REPR_FORMATr$_NON_OK_RENDEZVOUS_REPR_FORMATr.)rrr r r!_rpc_state_stringbs  rc@sVeZdZUdZeed<edddZeedddZ eedd d Z ee j dd d Z eedd dZeedddZedddZedddZedddZedddZedddZedddZedddZd*eeed d!d"Zd+eeeed d#d$Zd,eeeejd d%d&Z d-e!e j"gdfeedd'd(d)Z#dS)._InactiveRpcErrorzAn RPC error not tied to the execution of a particular RPC. The RPC represented by the state object must not be in-progress or cancelled. Attributes: _state: An instance of _RPCState. _stateruc Csv|j\tdt|jt|j|jt|j|_t|j |j_ t|j |j_ Wdn1sh0YdSrK) r(r'copydeepcopyr*r,r-r$rr+r.r;rJr r r!r<s   z_InactiveRpcError.__init__rcCs|jjSrrr*r=r r r!r*sz"_InactiveRpcError.initial_metadatacCs|jjSrrr,r=r r r!r,sz#_InactiveRpcError.trailing_metadatacCs|jjSrrr-r=r r r!r-sz_InactiveRpcError.codecCst|jjSr)r decoderr$r=r r r!r$sz_InactiveRpcError.detailscCst|jjSr)r rrr.r=r r r!r.sz$_InactiveRpcError.debug_error_stringcCst|jj|jSrr __class__r?rr=r r r!_reprsz_InactiveRpcError._reprcCs|Srrr=r r r!__repr__sz_InactiveRpcError.__repr__cCs|Srrr=r r r!__str__sz_InactiveRpcError.__str__cCsdS)zSee grpc.Future.cancel.Fr r=r r r!r~sz_InactiveRpcError.cancelcCsdS)zSee grpc.Future.cancelled.Fr r=r r r!r/sz_InactiveRpcError.cancelledcCsdS)zSee grpc.Future.running.Fr r=r r r!runningsz_InactiveRpcError.runningcCsdS)zSee grpc.Future.done.Tr r=r r r!risz_InactiveRpcError.doneNrcCs|dS)zSee grpc.Future.result.Nr r;rr r r!resultsz_InactiveRpcError.resultcCs|S)zSee grpc.Future.exception.r rr r r!r}sz_InactiveRpcError.exceptioncCs.z|Wn tjy(tdYS0dS)zSee grpc.Future.traceback.N)rDRpcErrorsysexc_inforr r r! tracebacksz_InactiveRpcError.traceback)fnrrcCs ||dS)z"See grpc.Future.add_done_callback.Nr )r;rrr r r!add_done_callbacksz#_InactiveRpcError.add_done_callback)N)N)N)N)$r?r@rA__doc__r'rBr<rrr*r,rDrEr-rFr$r.rrrrGr~r/rrirIrrrdr}types TracebackTyperrFuturerr r r r!rtsH      rcseZdZUdZeed<eejej fed<e e ed<e e ed<eeejej fe e e e dfdd Z ed d d Ze e d d d Zed ddZeedddZddZddZddZddZe ed ddZed ddZed dd Zed d!d"Zd#d d$d%ZZS)& _RendezvousaAn RPC iterator. Attributes: _state: An instance of _RPCState. _call: An instance of SegregatedCall or IntegratedCall. In either case, the _call object is expected to have operate, cancel, and next_event methods. _response_deserializer: A callable taking bytes and return a Python object. _deadline: A float representing the deadline of the RPC in seconds. Or possibly None, to represent an RPC with no deadline at all. r_call_response_deserializerr")rJrprNdeadlinecs*tt|||_||_||_||_dSr)superrr<rrrr")r;rJrprNrrr r!r<s z_Rendezvous.__init__rcCs8|jj|jjduWdS1s*0YdS)zSee grpc.RpcContext.is_activeNrr(r-r=r r r! is_actives z_Rendezvous.is_activecCsh|jjL|jdur$WddSt|jtdWdSWdn1sZ0YdS)z"See grpc.RpcContext.time_remainingNr)rr(r"maxrr=r r r!time_remainings  z_Rendezvous.time_remainingcCs|jj~|jjdurhtjj}d}|jtj ||d|j_ t |j|||jj WddSWddSWdn1s0YdS)zSee grpc.RpcContext.cancelNz!Locally cancelled by application!TF) rr(r-rDrE CANCELLEDrr~r rr/rLrc)r;r-r$r r r!r~s    z_Rendezvous.cancelrjrcCsf|jjJ|jjdur&WddS|jj|WddSWdn1sX0YdS)z See grpc.RpcContext.add_callbackNFT)rr(r0appendr;rjr r r! add_callbacks   z_Rendezvous.add_callbackcCs|Srr r=r r r!__iter__sz_Rendezvous.__iter__cCs|Sr_nextr=r r r!rysz_Rendezvous.nextcCs|Srrr=r r r!__next__sz_Rendezvous.__next__cCs tdSrNotImplementedErrorr=r r r!r!sz_Rendezvous._nextcCs tdSrrr=r r r!r.$sz_Rendezvous.debug_error_stringcCst|jj|jSrrr=r r r!r'sz_Rendezvous._reprcCs|Srrr=r r r!r*sz_Rendezvous.__repr__cCs|Srrr=r r r!r-sz_Rendezvous.__str__NcCs||jj`|jjdurZtjj|j_d|j_d|j_|j t j |jj|jj|jj Wdn1sn0YdS)Nz"Cancelled upon garbage collection!T) rr(r-rDrErr$r/rr~r rrcr=r r r!__del__0s    z_Rendezvous.__del__)r?r@rArr'rBr rSegregatedCallIntegratedCallrrrIr<rGrrr~rrrryrrrFr.rrrr __classcell__r r rr!rs.      rc@sFeZdZUdZeed<edddZedddZeddd Z edd d Z d'e e e d ddZd(e e e ed ddZd)e e e ejd ddZeejgd fd dddZe edddZe edddZe ejdddZe edddZe ej ddd Z!e dd!d"Z"e dd#d$Z#e edd%d&Z$d S)*_SingleThreadedRendezvousaNAn RPC iterator operating entirely on a single thread. The __next__ method of _SingleThreadedRendezvous does not depend on the existence of any other thread, including the "channel spin thread". However, this means that its interface is entirely synchronous. So this class cannot completely fulfill the grpc.Future interface. The result, exception, and traceback methods will never block and will instead raise an exception if calling the method would result in blocking. This means that these methods are safe to call from add_done_callback handlers. rrcCs |jjduSrrr=r r r! _is_completeOsz&_SingleThreadedRendezvous._is_completecCs4|jj|jjWdS1s&0YdSrrr(r/r=r r r!r/Rs z#_SingleThreadedRendezvous.cancelledcCs8|jj|jjduWdS1s*0YdSrrr=r r r!rVs z!_SingleThreadedRendezvous.runningcCs8|jj|jjduWdS1s*0YdSrrr=r r r!riZs z_SingleThreadedRendezvous.doneNrcCs~~|jj`|s tjd|jjtjjurF|jj WdS|jj rXt n|Wdn1sp0YdS)a9Returns the result of the computation or raises its exception. This method will never block. Instead, it will raise an exception if calling this method would otherwise result in blocking. Since this method will never block, any `timeout` argument passed will be ignored. zJ_SingleThreadedRendezvous only supports result() when the RPC is complete.N) rr(rrD experimental UsageErrorr-rErr+r/FutureCancelledErrorrr r r!r^s   z _SingleThreadedRendezvous.resultcCs~|jjh|s tjd|jjtjjur@WddS|jj rRt n|WdSWdn1sx0YdS)a*Return the exception raised by the computation. This method will never block. Instead, it will raise an exception if calling this method would otherwise result in blocking. Since this method will never block, any `timeout` argument passed will be ignored. zM_SingleThreadedRendezvous only supports exception() when the RPC is complete.N) rr(rrDrrr-rErr/rrr r r!r}us   z#_SingleThreadedRendezvous.exceptionc Cs~|jj|s tjd|jjtjjur@WddS|jj rRt n8z|Wn.tj yt dYWdS0Wdn1s0YdS)a;Access the traceback of the exception raised by the computation. This method will never block. Instead, it will raise an exception if calling this method would otherwise result in blocking. Since this method will never block, any `timeout` argument passed will be ignored. zM_SingleThreadedRendezvous only supports traceback() when the RPC is complete.Nr)rr(rrDrrr-rErr/rrrrrr r r!rs   z#_SingleThreadedRendezvous.tracebackrrcCsf|jjB|jjdur<|jjt||WddSWdn1sP0Y||dSrrr(r-r0rrrr;rr r r!rs   .z+_SingleThreadedRendezvous.add_done_callbackcCsJ|jj.|jjdur |q |jjWdS1s<0YdS)See grpc.Call.initial_metadataN)rr(r*_consume_next_eventr=r r r!r*s   z*_SingleThreadedRendezvous.initial_metadatacCsL|jj0|jjdur"tjd|jjWdS1s>0YdS)See grpc.Call.trailing_metadataNz4Cannot get trailing metadata until RPC is completed.)rr(r,rDrrr=r r r!r,s   z+_SingleThreadedRendezvous.trailing_metadatacCsL|jj0|jjdur"tjd|jjWdS1s>0YdS)See grpc.Call.codeNz'Cannot get code until RPC is completed.)rr(r-rDrrr=r r r!r-s   z_SingleThreadedRendezvous.codecCsR|jj6|jjdur"tjdt|jjWdS1sD0YdS)See grpc.Call.detailsNz*Cannot get details until RPC is completed.)rr(r$rDrrr rr=r r r!r$s   z!_SingleThreadedRendezvous.detailscCsV|j}|jj0t||j|j}|D] }|q(Wdn1sH0Y|Sr)r next_eventrr(rbr)r;rMr0rjr r r!rs   &z-_SingleThreadedRendezvous._consume_next_eventcCs||jjv|jjdur@|jj}d|j_|WdStjj|jjvrx|jjt j j urht n|jjdurx|Wdq1s0YqdSr) rrr(r+rrCrSr)r-rDrErrz)r;r+r r r!_next_responses   z(_SingleThreadedRendezvous._next_responsecCs|jjx|jjdurV|jjtjj|j t t fd}|sr|jj tjjn|jjt jjurntn|Wdn1s0Y|Sr)rr(r-r)rrrCrSrrReceiveMessageOperationrrQrDrErrzr)r;rr r r!rs   "z_SingleThreadedRendezvous._nextcCsR|jj6|jjdur"tjdt|jjWdS1sD0YdS)Nz5Cannot get debug error string until RPC is completed.)rr(r.rDrrr rr=r r r!r. s   z,_SingleThreadedRendezvous.debug_error_string)N)N)N)%r?r@rArr'rBrGrr/rrirrIrrrdr}rrrrrDrrrr*r,rEr-rFr$r BaseEventrrrr.r r r r!r=s,        rc@s$eZdZUdZeed<eedddZeedddZ ee j ddd Z ee dd d Zee dd d ZedddZedddZedddZedddZd#eeedddZd$eeeedddZd%eeeejdddZee jgdfdddd Zedd!d"Z dS)&_MultiThreadedRendezvousaAn RPC iterator that depends on a channel spin thread. This iterator relies upon a per-channel thread running in the background, dequeueing events from the completion queue, and notifying threads waiting on the threading.Condition object in the _RPCState object. This extra thread allows _MultiThreadedRendezvous to fulfill the grpc.Future interface and to mediate a bidirection streaming RPC. rrcsRjj6fdd}tjjj|jjWdS1sD0YdS)rcs jjduSrrr r=r r!rv&sz8_MultiThreadedRendezvous.initial_metadata.._doneN)rr(r rr*r;rvr r=r!r*"s  z)_MultiThreadedRendezvous.initial_metadatacsRjj6fdd}tjjj|jjWdS1sD0YdS)rcs jjduSrrr r=r r!rv0sz9_MultiThreadedRendezvous.trailing_metadata.._doneN)rr(r rr,rr r=r!r,,s  z*_MultiThreadedRendezvous.trailing_metadatacsRjj6fdd}tjjj|jjWdS1sD0YdS)rcs jjduSrrr r=r r!rv:sz,_MultiThreadedRendezvous.code.._doneN)rr(r rr-rr r=r!r-6s  z_MultiThreadedRendezvous.codecsXjj<fdd}tjjj|tjjWdS1sJ0YdS)rcs jjduSr)rr$r r=r r!rvDsz/_MultiThreadedRendezvous.details.._doneN)rr(r rrr$rr r=r!r$@s  z _MultiThreadedRendezvous.detailscsXjj<fdd}tjjj|tjjWdS1sJ0YdS)Ncs jjduSr)rr.r r=r r!rvMsz:_MultiThreadedRendezvous.debug_error_string.._done)rr(r rrr.rr r=r!r.Js  z+_MultiThreadedRendezvous.debug_error_stringcCs4|jj|jjWdS1s&0YdSrrr=r r r!r/Ss z"_MultiThreadedRendezvous.cancelledcCs8|jj|jjduWdS1s*0YdSrrr=r r r!rWs z _MultiThreadedRendezvous.runningcCs8|jj|jjduWdS1s*0YdSrrr=r r r!ri[s z_MultiThreadedRendezvous.donecCs |jjduSrrr=r r r!r_sz%_MultiThreadedRendezvous._is_completeNrcCs|jjrtj|jjj|j|d}|r0tn<|jjtjj urV|jj WdS|jj rht n|Wdn1s0YdS)zReturns the result of the computation or raises its exception. See grpc.Future.result for the full API contract. rN) rr(r rrrDFutureTimeoutErrorr-rErr+r/rr;r timed_outr r r!rbs   z_MultiThreadedRendezvous.resultcCs|jjztj|jjj|j|d}|r0tnD|jjtjj urPWddS|jj rbt n|WdSWdn1s0YdS)zvReturn the exception raised by the computation. See grpc.Future.exception for the full API contract. rN) rr(r rrrDrr-rErr/rrr r r!r}us   z"_MultiThreadedRendezvous.exceptionc Cs|jjtj|jjj|j|d}|r0tnj|jjtjj urPWddS|jj rbt n8z|Wn.tj yt dYWdS0Wdn1s0YdS)zAccess the traceback of the exception raised by the computation. See grpc.future.traceback for the full API contract. rNr)rr(r rrrDrr-rErr/rrrrrr r r!rs   z"_MultiThreadedRendezvous.tracebackrcCsf|jjB|jjdur<|jjt||WddSWdn1sP0Y||dSrrrr r r!rs   .z*_MultiThreadedRendezvous.add_done_callbackcs.jjjjdurftjj}jjtjj j t t f|}|sjjtjj njjtjjur~tnfdd}tjjj|jjdurΈjj}dj_|WdStjj jjvr jjtjjurtnjjdur Wdn1s 0YdS)Ncs(jjdup&tjjjjvo&jjduSr)rr+rrCrSr)r-r r=r r!_response_readys  z7_MultiThreadedRendezvous._next.._response_ready)rr(r-rnrr)rrrCrSrrrrrQrDrErrzr rr+)r;rrrrr+r r=r!rs4     z_MultiThreadedRendezvous._next)N)N)N)!r?r@rArr'rBrrr*r,rDrEr-rFr$r.rGr/rrirrIrrrdr}rrrrrrrr r r r!rs(        r)rrrqrcCsPt|}t||}|durBtdddtjjd}t|}|d|fS||dfSdS)Nr rs)r"r rr'rDrErVr)rrrqrrrJrfr r r!_start_unary_requests  r)rJrp with_callrrcCs>|jtjjur2|r*t||d|}|j|fS|jSnt|dSr)r-rDrErrr+r)rJrprr rendezvousr r r!_end_unary_response_blockings  r)metadatainitial_metadata_flagsrcCs*t||ttttfttffSr)rSendInitialMetadataOperationrrReceiveStatusOnClientOperationReceiveInitialMetadataOperationrrr r r!#_stream_unary_invocation_operationss rcCstddt||DS)Ncss|]}|dfVqdSrr ).0rr r r! sz?_stream_unary_invocation_operations_and_tags..)tuplerrr r r!,_stream_unary_invocation_operations_and_tagss r) user_deadlinercCsRt}|dur|durdS|dur0|dur0|S|durD|durD|St||SdSr)rget_deadline_from_contextmin)rparent_deadliner r r!_determine_deadlinesrc @seZdZUejed<eed<eed<eed<ee ed<ee ed<e ed<ee ed<gd Z ejeeeee ee ee d d d Ze eeeeeeeejeeeeeejeeeejfd ddZde eeeeeejeeeejeeejfdddZde eeeeeejeeeeje dddZde eeeeeejeeeejee ejfdddZde eeeeeejeeeeje dddZ!dS)_UnaryUnaryMultiCallable_channel _managed_call_method_target_request_serializerr_context_registered_call_handlerrrrrrrchannel managed_callr4r5rqrNr cCs8||_||_||_||_||_||_t|_||_ dSr rrrrrrrbuild_census_contextrr r;r r r4r5rqrNr r r r!r</s  z!_UnaryUnaryMultiCallable.__init__)rrrwait_for_ready compressionrc Cst|||j\}}}t|} t||} |dur@ddd|fSttdddd} t | | t |t t t t t tt tt f} | | |dfSdSr)rr_InitialMetadataFlagswith_wait_for_readyr augment_metadatar'_UNARY_UNARY_INITIAL_DUErrrrrrrr) r;rrrrrrrrraugmented_metadatarJrr r r!_prepareBs,     z!_UnaryUnaryMultiCallable._prepareNrrr credentialsrrrc Cs||||||\}}} } |dur(| nt|_t|j|_t|j|_ |j t j j|jdt| ||durtdn|j|dff|j|j } | } t| ||j|| fSdSr)rrr\r2r rrr4rr5rsegregated_callrPropagationConstantsGRPC_PROPAGATE_DEFAULTSr _credentialsrr rrbr) r;rrrrrrrJrrrrprMr r r! _blockinghs2   z"_UnaryUnaryMultiCallable._blockingc Cs&|||||||\}}t||ddSr6rr r;rrrrrrrJrpr r r!__call__s  z!_UnaryUnaryMultiCallable.__call__c Cs&|||||||\}}t||ddSNTr r!r r r!rs  z"_UnaryUnaryMultiCallable.with_callc Cs||||||\}}} } |dur(| nxt||j} t|_t|j|_ t|j |_ | t jj|jd| ||durzdn|j|f| |j|j } t|| |j| SdSr)rrnrrr\r2r rrr4rr5rrrrrrr r) r;rrrrrrrJrrrrrrpr r r!futures0      z_UnaryUnaryMultiCallable.future)NNNNN)NNNNN)NNNNN)NNNNN)"r?r@rArChannelrBrbytesrrrrrH __slots__r<rIrrGrD Compressionr r'r OperationrrCallCredentialsrrr"Callrrr$r r r r!rs         ) )  rc @seZdZUejed<eed<eed<eeed<ee ed<e ed<ee ed<gdZ ejeeee ee d d d Z de eeeeeejeeeejed ddZd S)'_SingleThreadedUnaryStreamMultiCallablerrrrrrr )rrrrrr)r r4r5rqrNr cCs2||_||_||_||_||_t|_||_dSr) rrrrrrrrr )r;r r4r5rqrNr r r r!r<s  z0_SingleThreadedUnaryStreamMultiCallable.__init__Nrc Cst|}t||j}|dur:tdddtjjd} t| tt dddd} |durVdn|j } t |} t ||} t| | t|tttfttfttff} tdd| D}t| _t|j| _t|j| _|j tj!j"|jdt#||| ||j$|j% }t&| ||j'|S)Nr rscss|]}|dfVqdSrr )ropsr r r!r$zC_SingleThreadedUnaryStreamMultiCallable.__call__..)(r"r rrr'rDrErVr_UNARY_STREAM_INITIAL_DUErrrr rrrrrrrrrrr\r2rrr4rr5rrrrrrr rr)r;rrrrrrrrrJcall_credentialsrrroperations_and_tagsrpr r r!r"sb        z0_SingleThreadedUnaryStreamMultiCallable.__call__)NNNNN)r?r@rArr%rBr&rrrrrHr'r<rIrrDr*rGr(rr"r r r r!r,s:       r,c @seZdZUejed<eed<eed<eed<ee ed<ee ed<e ed<ee ed<gd Z ejeeee e ee d d d Zde eeeeeejeeeejedddZd S)_UnaryStreamMultiCallablerrrrrrrr r r cCs8||_||_||_||_||_||_t|_||_ dSrrrr r r!r<Ms  z"_UnaryStreamMultiCallable.__init__Nrc Cst|||j\}}} t|} |dur.| nt||} ttdddd} t | | t |t t t t t ftt ff} t| _t|j| _t|j| _|tjj|jdt|||durdn|j| t| |j|j|j }t!| ||j|SdSr)"rrrrr rr'r/rrrrrrrrr\r2r rrr4rr5rrrrrrnrrr r)r;rrrrrrrrrrrrJrrpr r r!r"`sR       z"_UnaryStreamMultiCallable.__call__)NNNNN)r?r@rArr%rBrr&rrrrrHr'r<rIrrDr*rGr(rr"r r r r!r28s>       r2c @sneZdZUejed<eed<eed<eed<ee ed<ee ed<e ed<ee ed<gd Z ejeeeee ee ee d d d Zeeeeeeejeeeejeeejfd ddZdeeeeeeejeeeeje d ddZdeeeeeeejeeeejee ejfd ddZdeeeeeeejeeeejed ddZdS)_StreamUnaryMultiCallablerrrrrrrr r r cCs8||_||_||_||_||_||_t|_||_ dSrrrr r r!r<s  z"_StreamUnaryMultiCallable.__init__rorrrrrrc Cs t|}ttdddd}t|} t||} t|_ t |j |_ t |j|_|jtjj|j dt|| |dur|dn|jt| | |j|j } t||| |jd| } |j>t| ||j|j |j!sWdqWdq1s0Yq|| fSr)"r"r'_STREAM_UNARY_INITIAL_DUErrr rrr\r2r rrr4rr5rrrrrrrrrr rrrr(rbrrcr)) r;rorrrrrrrJrrrprMr r r!rsD     0z#_StreamUnaryMultiCallable._blockingNc Cs&|||||||\}}t||ddSr6r  r;rorrrrrrJrpr r r!r"s  z"_StreamUnaryMultiCallable.__call__c Cs&|||||||\}}t||ddSr#r r6r r r!rs  z#_StreamUnaryMultiCallable.with_callc Cst|}ttdddd}t||j} t|} t||} t |_ t |j|_t |j|_|tjj|jd|| |durdn|jt|| | |j|j } t||| |j| t|| |j|Sr)r"r'r5rnrrrr rrr\r2r rrr4rr5rrrrrrrr rrr) r;rorrrrrrrJrrrrrpr r r!r$sH    z _StreamUnaryMultiCallable.future)NNNNN)NNNNN)NNNNN)r?r@rArr%rBrr&rrrrrHr'r<rrIrrDr*rGr(r r'rrr"r+rrr$r r r r!r3s        0  r3c @seZdZUejed<eed<eed<eed<ee ed<ee ed<e ed<ee ed<gd Z ejeeeee ee ee d d d ZdeeeeeeejeeeejedddZd S)_StreamStreamMultiCallablerrrrrrrr r r cCs8||_||_||_||_||_||_t|_||_ dSrrrr r r!r<\s  z#_StreamStreamMultiCallable.__init__Nr4c Cst|}ttdddd}t|} t||} t| | t t ft t ff} t ||j } t|_t|j|_t|j|_|tjj|jdt|| |durdn|j| | |j|j } t||| |j| t || |j |Sr)!r"r'_STREAM_STREAM_INITIAL_DUErrr rrrrrrrnrrr\r2r rrr4rr5rrrrrrr rrr)r;rorrrrrrrJrrrrrrpr r r!r"osR      z#_StreamStreamMultiCallable.__call__)NNNNN)r?r@rArr%rBrr&rrrrrHr'r<rrIrrDr*rGr(rr"r r r r!r7Gs>       r7cs>eZdZdZefedfdd ZeeedddZ Z S)rz'Stores immutable initial metadata flags)valuecs|tjjM}tt|||Sr)rInitialMetadataFlags used_maskrr__new__)clsr9rr r!r<s z_InitialMetadataFlags.__new__)rrcCsJ|durF|r&||tjjBtjjBS|sF||tjj@tjjBS|Sr)rrr:rwait_for_ready_explicitly_set)r;rr r r!rs  z)_InitialMetadataFlags.with_wait_for_ready) r?r@rArrrHr<rrGrrr r rr!rsrc@sNeZdZUejed<eed<eed<ejdddZddd d Z d d Z dS) _ChannelCallStater  managed_callsr7r cCs t|_||_d|_d|_dS)NrF)r7Locklockr r@r;r r r r!r<s z_ChannelCallState.__init__NrcCs d|_dS)Nr)r@r=r r r!r>sz&_ChannelCallState.reset_postfork_childc Cs2z|jtjjdWnttfy,Yn0dS)NzChannel deallocated!)r closerrEr/ TypeErrorAttributeErrorr=r r r!rs z_ChannelCallState.__del__) r?r@rArr%rBrHrGr<r>rr r r r!r?s  r?)rJrcs.fdd}tj|d}|d|dS)Ncstj}|jtjjkr$q||}|rj8j d8_ j dkrbWddSWdq1sv0YqdS)Nr) rrr next_call_eventcompletion_typeCompletionType queue_timeouttagrCr@)rMcall_completedrur r! channel_spins    z._run_channel_spin_thread..channel_spinrTr)rJrOchannel_spin_threadr rur!_run_channel_spin_threads  rQruc sLttttttttttjtttj t t tttj d fdd }|S)N) flagsr4hostrrrrrrcontextr rc stfdd|D} jXj||||||| || } jdkrTd_tnjd7_| WdS1sz0YdS)aCreates a cygrpc.IntegratedCall. Args: flags: An integer bitfield of call flags. method: The RPC method. host: A host string for the created call. deadline: A float to be the deadline of the created call or None if the call is to have an infinite deadline. metadata: The metadata for the call or None. credentials: A cygrpc.CallCredentials or None. operations: A sequence of sequences of cygrpc.Operations to be started on the call. event_handler: A behavior to call to handle the events resultant from the operations on the call. context: Context object for distributed tracing. _registered_call_handle: An int representing the call handle of the method, or None if the method is not registered. Returns: A cygrpc.IntegratedCall with which to conduct an RPC. c3s|]}|fVqdSrr )r operationrrr r!rszC_channel_managed_call_management..create..rrHN)rrCr integrated_callr@rQ) rRr4rSrrrrrrrTr r1rprurVr!creates(   z0_channel_managed_call_management..create) rHr&rrFrIrrr*rr)rrr)rJrXr rur! _channel_managed_call_managements :rYc@seZdZUejed<ejed<eed<ej ed<eed<e e e e ej gdfeej fed<eed<ejd d d Zdd d dZdS)_ChannelConnectivityStaterCr polling connectivitytry_to_connectNcallbacks_and_connectivities deliveringrAcCs2t|_||_d|_d|_d|_g|_d|_dSr6) r7RLockrCr r[r\r]r^r_rDr r r!r<6s z"_ChannelConnectivityState.__init__rcCs"d|_d|_d|_g|_d|_dSr6)r[r\r]r^r_r=r r r!r>?s z._ChannelConnectivityState.reset_postfork_child)r?r@rAr7r`rBrDr%rGChannelConnectivityrrr rrr<r>r r r r!rZ%s"     rZcCs:g}|jD]*}|\}}||jur |||j|d<q |S)NrH)r^r\r)rJcallbacks_needing_updatecallback_and_connectivityrjcallback_connectivityr r r! _deliveriesGs    re)rJinitial_connectivityinitial_callbacksrc Cs|}|}|D]8}t|z ||Wq tyBttYq 0q |j:t|}|rb|j}nd|_ WddSWdq1s0YqdSr6) rrrdr|r}0_CHANNEL_SUBSCRIPTION_CALLBACK_ERROR_LOG_MESSAGErCrer\r_)rJrfrgr\r0rjr r r!_deliverVs     ri)rJr0rcCs2tjt||j|fd}|d|d|_dSN)r5argsT)rrrir\rrr_)rJr0delivering_threadr r r!_spawn_deliveryos rm)rJr initial_try_to_connectrcCs^|}||}|jTtj||_tdd|jD}|jD]}|j|d<q<|rZt||Wdn1sn0Y||t d}t ||jD|js|j sd|_ d|_WdqZ|j }d|_ Wdn1s0Y|js|rx||}|j<tj||_|js8t|}|r8t||Wdqx1sN0YqxdS)Ncss|]\}}|VqdSrr )rrj_r r r!rsz%_poll_connectivity..rHg?F)check_connectivity_staterCr 1CYGRPC_CONNECTIVITY_STATE_TO_CHANNEL_CONNECTIVITYr\rr^rmwatch_connectivity_staterrrr]r[successr_re)rJr rnr]r\r0rcrMr r r!_poll_connectivitysN   (  $  rt)rJrjr]rcCs|j|jsX|jsXtjt||jt|fd}|d| d|_|j |dgnd|j s|j durt ||f|jt|O_|j ||j gn"|jt|O_|j |dgWdn1s0YdSrj)rCr^r[rrrtr rGrrrr_r\rmr])rJrjr]polling_threadr r r! _subscribes$   rv)rJrjrcCsZ|j@t|jD]$\}\}}||kr|j|q8qWdn1sL0YdSr)rC enumerater^pop)rJrjindexsubscribed_callbackunused_connectivityr r r! _unsubscribes r|) base_optionsrrcCs$t|}t||tjjtffSr)r create_channel_optionrr ChannelArgKeyprimary_user_agent_string _USER_AGENT)r}rcompression_optionr r r!_augment_optionss r)optionsrcCsBg}g}|D],}|dtjjjkr.||q ||q ||fS)z;Separates core channel options from Python channel options.r)rDrChannelOptionsSingleThreadedUnaryStreamr)r core_optionspython_optionspairr r r!_separate_channel_optionss  rc@seZdZUdZeed<ejed<eed<e ed<e ed<e e e fed<e e eeejeejdd d Ze e d d d Ze eddddZd1eejgdfeeddddZeejgdfddddZd2e eeeeeeejdddZd3e eeeeeeejdddZd4e eeeeeeej dddZ!d5e eeeeeeej"dd d!Z#dd"d#d$Z$dd"d%d&Z%dd"d'd(Z&d)d*Z'd+d,Z(dd"d-d.Z)d/d0Z*dS)6r%z7A cygrpc.Channel-backed implementation of grpc.Channel._single_threaded_unary_streamr _call_state_connectivity_stater_registered_call_handles)r5rrrcCsrt|\}}t|_||tt|t||||_ ||_ t |j |_ t |j |_t|tjrntdS)aPConstructor. Args: target: The target to which to connect. options: Configuration options for the channel. credentials: A cygrpc.ChannelCredentials or None. compression: An optional value indicating the compression method to be used over the lifetime of the channel. N)r%_DEFAULT_SINGLE_THREADED_UNARY_STREAMr_process_python_optionsrr%r encoderrrr?rrZrfork_register_channelg_gevent_activatedgevent_increment_channel_count)r;r5rrrrrr r r!r<s     zChannel.__init__)r4rcCs|jt|S)ah Get the registered call handle for a method. This is a semi-private method. It is intended for use only by gRPC generated code. This method is not thread-safe. Args: method: Required, the method name for the RPC. Returns: The registered call handle pointer in the form of a Python Long. )rget_registered_call_handler r)r;r4r r r!_get_registered_call_handle&sz#Channel._get_registered_call_handleN)rrcCs&|D]}|dtjjjkrd|_qdS)zASets channel attributes according to python-only channel options.rTN)rDrrrr)r;rrr r r!r6s zChannel._process_python_options)rjr]rcCst|j||dSr)rvr)r;rjr]r r r! subscribeAszChannel.subscribercCst|j|dSr)r|rrr r r! unsubscribeHszChannel.unsubscribeF)r4rqrN_registered_methodrcCs<d}|r||}t|jt|jt|t|j|||Sr)rrrrYrr rrr;r4rqrNrr r r r! unary_unaryNs  zChannel.unary_unarycCshd}|r||}|jr:t|jt|t|j|||St|jt|j t|t|j|||SdSr) rrr,rr rrr2rYrrr r r! unary_streamcs*    zChannel.unary_streamcCs<d}|r||}t|jt|jt|t|j|||Sr)rr3rrYrr rrrr r r! stream_unarys  zChannel.stream_unarycCs<d}|r||}t|jt|jt|t|j|||Sr)rr7rrYrr rrrr r r! stream_streams  zChannel.stream_streamrcCs@|j}|r<|j|jdd=Wdn1s20YdSr)rrCr^rr r r!_unsubscribe_allszChannel._unsubscribe_allcCs6||jtjjdt|tjr2tdS)NzChannel closed!) rrrErrEr/fork_unregister_channelrgevent_decrement_channel_countr=r r r!_closes  zChannel._closecCs||jtjjddS)NzChannel closed due to fork)rr close_on_forkrrEr/r=r r r!_close_on_forkszChannel._close_on_forkcCs|Srr r=r r r! __enter__szChannel.__enter__cCs |dSr6r)r;exc_typeexc_valexc_tbr r r!__exit__szChannel.__exit__cCs |dSrrr=r r r!rEsz Channel.closecCsz |Wn Yn0dSr)rr=r r r!rs  zChannel.__del__)N)NNF)NNF)NNF)NNF)+r?r@rArrGrBrr%r?rZrFrrHrrrrDChannelCredentialsr(r<rrrrarrrrUnaryUnaryMultiCallablerUnaryStreamMultiCallablerStreamUnaryMultiCallablerStreamStreamMultiCallablerrrrrrrErr r r r!r%s   !     &  r%)urrrreosrr7rrtypingrrrrrrrr r r rDr r rr grpc._cythonr grpc._typingrrrrrrrrgrpc.experimental getLoggerr?r|r% __version__rrgetenvrrCsend_initial_metadatartrrRrSrWrr/r5r8rhrrrIr"rErFr&objectr'rLrrbrnrrrrrr+rr RpcContextrrrr&rrGrrHr)rrrrrrr,r2rr3rr7rr?rQrYrZrarerirmr%rtrvr|r(rrr r r r!s20                    =  -  _^k  W  ;     ;d_1^?#    2