5,j[dUdZddlmZddlZddlZddlmZddlmZddl m Z m Z ddl Z ddl mZmZmZmZe rddlmZd d lmZd d lmZ d Zd ZdZdZdZGddeZGddeZGddeZ dCdZ!Gddej"Z#GddZ$Gd d!ej"Z%dDd#Z&dDd$Z'dDd%Z(dDd&Z)dZ*d'e+d(<e)r)e're&se(s e$Z*ndZ*d)Z,dEd-Z-d.a.dFd0Z/d.a0dFd1Z1dGd3Z2dHd5Z3dId=Z4dDd>Z5dDd?Z6dJdAZ7dDdBZ8dS)Ka Support for streaming http requests in emscripten. A few caveats - If your browser (or Node.js) has WebAssembly JavaScript Promise Integration enabled https://github.com/WebAssembly/js-promise-integration/blob/main/proposals/js-promise-integration/Overview.md *and* you launch pyodide using `pyodide.runPythonAsync`, this will fetch data using the JavaScript asynchronous fetch api (wrapped via `pyodide.ffi.call_sync`). In this case timeouts and streaming should just work. Otherwise, it uses a combination of XMLHttpRequest and a web-worker for streaming. This approach has several caveats: Firstly, you can't do streaming http in the main UI thread, because atomics.wait isn't allowed. Streaming only works if you're running pyodide in a web worker. Secondly, this uses an extra web worker and SharedArrayBuffer to do the asynchronous fetch operation, so it requires that you have crossOriginIsolation enabled, by serving over https (or from localhost) with the two headers below set: Cross-Origin-Opener-Policy: same-origin Cross-Origin-Embedder-Policy: require-corp You can tell if cross origin isolation is successfully enabled by looking at the global crossOriginIsolated variable in JavaScript console. If it isn't, streaming requests will fallback to XMLHttpRequest, i.e. getting the whole request into a buffer and then returning it. it shows a warning in the JavaScript console in this case. Finally, the webworker which does the streaming fetch is created on initial import, but will only be started once control is returned to javascript. Call `await wait_for_streaming_ready()` to wait for streaming fetch. NB: in this code, there are a lot of JavaScript objects. They are named js_* to make it clear what type of object they are. ) annotationsN)Parser)files) TYPE_CHECKINGAny)JsArray JsExceptionJsProxyto_js)Buffer)EmscriptenRequest)EmscriptenResponse)z user-agentc,eZdZ d dddd fd ZxZS) _RequestErrorNrequestresponsemessage str | NonerEmscriptenRequest | NonerEmscriptenResponse | Nonec~||_||_||_t|jdSN)rrrsuper__init__)selfrrr __class__s PC:\PYTHON\GemmaClient\venv\Lib\site-packages\urllib3/contrib/emscripten/fetch.pyr z_RequestError.__init__Hs:     &&&&&r)rrrrrr)__name__ __module__ __qualname__r __classcell__r"s@r#rrGsY# '-1.2 ' ' ' ' ' ' ' ' ' ' ' 'r$rceZdZdS)_StreamingErrorNr%r&r'r$r#r+r+UDr$r+ceZdZdS) _TimeoutErrorNr,r-r$r#r0r0Yr.r$r0dict_valdict[str, Any]returnr cBt|tjjS)N)dict_converter)r jsObject fromEntries)r1s r#_obj_from_dictr9]s ")*? @ @ @@r$cpeZdZdd ZddZddZeddZdfd ZddZ ddZ ddZ ddZ xZ S) _ReadStream int_bufferr byte_buffertimeoutfloatworkerr connection_idintrrc||_||_d|_d|_||_||_|dkrt d|znd|_d|_d|_ ||_ dS)NrTF) r<r=read_posread_lenrAr@rBr>is_live _is_closedr)r!r<r=r>r@rArs r#r z_ReadStream.__init__bsi%&  * .5kks4'>***t  18 r$r3Nonec.|dSrcloser!s r#__del__z_ReadStream.__del__v r$boolc|jSrrHrMs r# is_closedz_ReadStream.is_closedz r$c*|SrrSrMs r#closedz_ReadStream.closed~~~r$c@|rdSd|_d|_d|_d|_d|_d|_|jr5|j td|j id|_t dS)NrTrLF)rSrFrEr<r=rHrrGr@ postMessager9rArrLr!r"s r#rLz_ReadStream.closes >>    F   < ! K # #NGT=O3P$Q$Q R R R DL  r$cdSNTr-rMs r#readablez_ReadStream.readabletr$cdSNFr-rMs r#writablez_ReadStream.writableur$cdSrar-rMs r#seekablez_ReadStream.seekablercr$byte_objr c8|jstd|jd|jdkrRtj|jdt|j td|j itj |jdt|j dkrt|jd}|dkr||_d|_n|t krs|jd}tj}||jd|}td||jdd|_|dSt1|jt3t5|}|j|j|j|z}|t5|d|<|xj|zc_|xj|z c_|S) Nz,No buffer for stream in _ReadStream.readintorrgetMorez timed-outr Exception thrown in fetch: F)r<r+rrFr6Atomicsstore ERROR_TIMEOUTr@rZr9rAwaitr>r0rEERROR_EXCEPTION TextDecodernewdecoder=slicerGrLminlen memoryviewsubarrayto_py)r!rfdata_len string_len js_decoderjson_str ret_lengthrvs r#readintoz_ReadStream.readintos !>   =A   J  T_a ? ? ? K # #NIt?Q3R$S$S T T T M4<PP$#q)H!|| ( ! _,,!_Q/ ^//11 %,,T-=-C-CAz-R-RSS%<(<< L! %  qJx,@,@(A(ABB #,, M4=:5  %'' .6 8Qz\* #  # r$) r<rr=rr>r?r@r rArBrrr3rIr3rPrfr r3rB)r%r&r'r rNrSpropertyrWrLr^rbrer}r(r)s@r#r;r;as9999(   X       ,,,,,,,,r$r;ceZdZd dZd dZdS) _StreamingFetcherr3rIcd_ttdd}t jt|gdtddi}dfd }t j |}t j j |_t j j|_dS)NFzemscripten_fetch_worker.jszutf-8)encoding)create_pyproxiestypezapplication/javascript js_resolve_fnr js_reject_fnr3rIcVdfd }dfd }|j_|j_dS)Ner r3rIc,d_|dSr])streaming_ready)rrr!s r#onMsgzC_StreamingFetcher.__init__..promise_resolver..onMsgs!'+$ a     r$c|dSrr-)rrs r#onErrzC_StreamingFetcher.__init__..promise_resolver..onErrs Qr$)rr r3rI) js_worker onmessageonerror)rrrrr!s`` r#promise_resolverz4_StreamingFetcher.__init__..promise_resolversb ! ! ! ! ! ! !      (-DN $%*DN " " "r$)rr rr r3rI)rr __package__joinpath read_textr6Blobrpr r9URLcreateObjectURL globalThisWorkerrPromisejs_worker_ready_promise)r!streaming_worker_code js_data_blobr js_data_urls` r#r z_StreamingFetcher.__init__s$ +   X2 3 3 YY ( (  w{{ ()E B B B F$<= > >  + + + + + +f,,\:: -11+>>')}'<'@'@AQ'R'R$$$r$rrrc d|jD}|j}|t||jd}|jdkrt d|jznd}tj d}tj |}tj |d}tj |dttj |dtj |jtjj} |jt-|| |dtj |dt||dtkrt1d|d |dt2kr|d } tj } | |d| } t;j| } t?|| d | d tA|||j|j| d |S|dtBkrd|d } tj } | |d| } tEd| |d tEd|d|d )Nc,i|]\}}|tv||Sr-HEADERS_TO_IGNORE.0kvs r# z*_StreamingFetcher.send..s0   QAR8R8RAq8R8R8Rr$)headersbodymethodrrDi)bufferurl fetchParamsz'Timeout connecting to streaming requestrr statusr connectionID)r status_coderrriz%Unknown status from worker in fetch: )#ritemsrr rr>rBr6SharedArrayBufferrp Int32Array Uint8ArrayrjrkrlnotifyrrlocationhrefrrZr9rmr0SUCCESS_HEADERrorqrrjsonloadsrr;rnr+)r!rrr fetch_datar>js_shared_buffer js_int_bufferjs_byte_bufferjs_absolute_urlryrzr{ response_objs r#sendz_StreamingFetcher.sends  $_2244   |!(%++XX 1811D1D#dW_,---$/33G<< ))*:;; **+;Q?? =999 -+++&**W["+>>C "" .*#-       q-AAA  } , ,9  1  / /'q)J++--J"(()=)=a)L)LMMH:h//L%(2$Y/ !"ON 0     1  0 0&q)J++--J!(()=)=a)L)LMMH!8h88'TX "J a8HJJ r$Nr~rrr3r)r%r&r'r rr-r$r#rrsFSSSS8FFFFFFr$rc|eZdZdZdd ZddZddZeddZdfd Z ddZ ddZ ddZ ddZ ddZxZS)_JSPIReadStreamaF A read stream that uses pyodide.ffi.run_sync to read from a JavaScript fetch response. This requires support for WebAssembly JavaScript Promise Integration in the containing browser, and for pyodide to be launched via runPythonAsync. :param js_read_stream: The JavaScript stream reader :param timeout: Timeout in seconds :param request: The request we're handling :param response: The response this stream relates to :param js_abort_controller: A JavaScript AbortController object, used for timeouts js_read_streamrr>r?rrrrjs_abort_controllerc||_||_d|_d|_||_||_d|_d|_||_dS)NFr) rr>rH_is_donerrcurrent_buffercurrent_buffer_posr)r!rr>rrrs r#r z_JSPIReadStream.__init__DsM-  18 3; ""##6   r$r3rIc.|dSrrKrMs r#rNz_JSPIReadStream.__del__VrOr$rPc|jSrrRrMs r#rSz_JSPIReadStream.is_closedZrTr$c*|SrrVrMs r#rWz_JSPIReadStream.closed^rXr$c|rdSd|_d|_|jd|_d|_d|_d|_d|_t dS)NrT) rSrFrErcancelrHrrrrrLr[s r#rLz_JSPIReadStream.closebsv >>    F   ""$$$"     r$cdSr]r-rMs r#r^z_JSPIReadStream.readableor_r$cdSrar-rMs r#rbz_JSPIReadStream.writablerrcr$cdSrar-rMs r#rez_JSPIReadStream.seekableurcr$ct|j|j|j|j|j}|jr d|_dS|j |_ d|_ dS)NrTFr) _run_sync_with_timeoutrreadr>rrrdonervaluerwrr)r! result_jss r#_get_next_bufferz _JSPIReadStream._get_next_bufferxsx*   $ $ & & L  $L]     >  DM5"+/"7"7"9"9D &'D #4r$rfr rBc|j1|r|j|dStt |t |j|jz }|j|j|j|z|d|<|xj|z c_|jt |jkrd|_|S)Nr)rrrLrsrtr)r!rfr|s r#r}z_JSPIReadStream.readintos   &((** d.A.I q MM3t233d6MM  "&!4  #d&= &J J" : :-  "c$*=&>&> > >"&D r$) rrr>r?rrrrrrr~rr)r%r&r'__doc__r rNrSrrWrLr^rbrerr}r(r)s@r#rr.s*7777$   X        r$rrPcttdo.ttdotjtjkS)Nwindowr!)hasattrr6r!rr-r$r#is_in_browser_main_threadrs/ 2x QWR%8%8 QRW =QQr$cDttdo tjS)NcrossOriginIsolated)rr6rr-r$r#is_cross_origin_isolatedrs 2, - - H"2HHr$cttdoRttjdo8ttjjdotjjjdkS)Nprocessreleasenamenoderr6rrrr-r$r# is_in_nodersVI . BJ * * . BJ& / / . J  #v - r$cVttdottdS)Nrr)rr6r-r$r#is_worker_availablers! 2x 8WR%8%88r$z_StreamingFetcher | None_fetcherzurllib3 only works in Node.js with pyodide.runPythonAsync and requires the flag --experimental-wasm-stack-switching in versions of node <24.rrrctrt|dStrtt|dt r(t rt |StdS)NTrrr) has_jspisend_jspi_requestrrNODE_JSPI_ERRORrrr_show_streaming_warningrs r#send_streaming_requestrszz  $///  #    O%%}}W%%%!!!tr$FrIc^ts%dad}tj|dSdS)NTz8Warning: Timeout is not available on main browser thread)_SHOWN_TIMEOUT_WARNINGr6consolewarn)rs r#_show_timeout_warningrs9 !!!%L      !!r$ctsodad}ts|dz }tr|dz }ts|dz }t dur|dz }dd lm}||dSdS) NTz%Can't stream HTTP requests because: z$ Page is not cross-origin isolated z+ Python is running in main browser thread z> Worker or Blob classes are not available in this environment.Fz Streaming fetch worker isn't ready. If you want to be sure that streaming fetch is working, you need to call: 'await urllib3.contrib.emscripten.fetch.wait_for_streaming_ready()`r)r)_SHOWN_STREAMING_WARNINGrrrrr6rr)rrs r#rrs ##' :')) ? > >G $ & & F E EG"$$ X W WG    % % e eG Wr$rctrt|dStrtt|d t j}ts+d|_ |j rt|j dz|_ n*| d|j rt||j|jd|jD]6\}}|t(vr|||7|t/|jt3t5|}ts,|j}n|j d}tC|j"|||S#tF$r]}|j$dkrtK|j&| |j$d krt|j&| t|j&| d}~wwxYw) NFr arraybufferrDztext/plain; charset=ISO-8859-15z ISO-8859-15rrrr TimeoutErrorr NetworkError)'rrrrrr6XMLHttpRequestrpr responseTyper>rBoverrideMimeTyperopenrrrrlowerrsetRequestHeaderrr rdictrparsestrgetAllResponseHeadersrrwtobytesencoderrr rr0r)rjs_xhrrrrrerrs r# send_requestrsVzz  %000  #    %>"&&(((** ("/F  =!$W_t%;!>> 8~ % % W=== = X ' ' W=== =  W=== =>sGH I1AI,,I1 streamingc|j}tj}d|jD}|j}|t||j|j d}trd|d<tj |j t|}t||||d}i}|j} | } t#| dd rn6t%| jd |t%| jd <\|j} d } t+| |d | } |r4|j,|j}t/|||| |} n8t||||| } | | _| S)a7 Send a request using WebAssembly JavaScript Promise Integration to wrap the asynchronous JavaScript fetch api (experimental). :param request: Request to send :param streaming: Whether to stream the response :return: The response object :rtype: EmscriptenResponse c,i|]\}}|tv||Sr-rrs r#rz%send_jspi_request..6s)VVV11DU;U;Uq!;U;U;Ur$)rrrsignalmanualredirectNrTrFr rr$r)r>r6AbortControllerrprrrr rr _is_node_jsfetchrr9rentriesnextgetattrstrrrr getReaderr arrayBufferrw)rrr>rrreq_bodyrfetcher_promise_js response_js header_iter iter_value_jsrrrbody_stream_jss r#rr$s oG,0022VV 5 5 7 7VVVG|Hh.%, J}}*!) :'+~j/I/IJJ) KG%--//KO#((** =&% 0 0 O 36}7J17M3N3NGC +A.// 0 O $K!$D!sGH   '(-7799N"(r?rrrcd}|dkr=tj|j|t |dz} ddlm}|||tj|SS#t$r9}|j dkrtd||t|j ||d}~wwxYw#|tj|wwxYw)ak Await a JavaScript promise synchronously with a timeout which is implemented via the AbortController :param promise: Javascript promise to await :param timeout: Timeout in seconds :param js_abort_controller: A JavaScript AbortController object, used on timeout :param request: The request being handled :param response: The response being handled (if it exists yet) :raises _TimeoutError: If the request times out :raises _RequestError: If the request raises a JavaScript exception :return: The result of awaiting the promise. NrrD)run_sync AbortErrorzRequest timed outr) r6 setTimeoutabortbindrB pyodide.ffir* clearTimeoutr rr0rr)r(r>rrrtimer_idr*rs r#rrrs>H{{=  % * *+> ? ?Wt^ATAT  &((((((x     OH % % % %  YYY 8| # #+Wx   WxXXX X Y   OH % % % % s$A// B294B--B22B55Ccd ddlm}m}t|S#t$rYdSwxYw)a Return true if jspi can be used. This requires both browser support and also WebAssembly to be in the correct state - i.e. that the javascript call into python was async not sync. :return: True if jspi can be used. :rtype: bool r can_run_syncr*F)r/r4r*rP ImportErrorr3s r#rrsU66666666LLNN### uus ! //cttdo3ttjdotjjjdkS)z_ Check if we are in Node.js. :return: True if we are in Node.js. :rtype: bool rrrrr-r$r#rrsA I . BJ * * . J  #v - r$ bool | Nonec,tr tjSdSr)rrr-r$r#rrs''tr$c@Ktrtjd{VdSdS)NTF)rrr-r$r#wait_for_streaming_readyr:s3........tur$)r1r2r3r r)rrr3rr~r)rrrrPr3r) r(rr>r?rrrrrrr3r)r3r7)9r __future__rior email.parserrimportlib.resourcesrtypingrrr6r/rr r r typing_extensionsr rrrrrr SUCCESS_EOFrlrn Exceptionrr+r0r9 RawIOBaser;rrrrrrr__annotations__rrrrrrrrrrrrr:r-r$r#rEs"""H#""""" %%%%%%%%%%%%%% )((((((&&&&&&(((((($   ' ' ' ' 'I ' ' '     m        M   AAAAddddd",dddNccccccccLhhhhhblhhhXRRRRIIII9999&*))))(A(A(C(C Z\\! ""HHH"!!!!!&.>.>.>.>bKKKK\3&3&3&3&l&    r$