Warning: session_start(): Session cannot be started after headers have already been sent in /home/tvrreohg/public_html/manga.php on line 13
§ Ìß]jÀ/ãóª—ddlmZddlZddlmZmZmZmZmZm Z ddl m Z ddl m Z ddlmZerddlmZdZd ZGd „d ¦«ZGd „d ¦«ZdS)é)Ú annotationsN)Ú TYPE_CHECKINGÚ AsyncIteratorÚ AwaitableÚCallableÚListÚOptional)Úuuid4)Úerrors)ÚMsg)ÚJetStreamContextiicóô—eZdZdZddddddeefd'd„Zed(d„¦«Zed(d„¦«Z ed)d„¦«Z ed*d„¦«Z ed*d„¦«Z ed*d„¦«Z d+d,d„Zd „Zd!„Zd-d"„Zd.d/d$„Zd-d%„Zd-d&„ZdS)0Ú Subscriptiona÷ A Subscription represents interest in a particular subject. A Subscription should not be constructed directly, rather `connection.subscribe()` should be used to get a subscription. :: nc = await nats.connect() # Async Subscription async def cb(msg): print('Received', msg) await nc.subscribe('foo', cb=cb) # Sync Subscription sub = nc.subscribe('foo') msg = await sub.next_msg() print('Received', msg) rÚNÚidÚintÚsubjectÚstrÚqueueÚcbú*Optional[Callable[[Msg], Awaitable[None]]]ÚfutureúOptional[asyncio.Future]Úmax_msgsÚpending_msgs_limitÚpending_bytes_limitÚreturnÚNonec ó.—||_||_||_||_||_d|_||_||_d|_||_ | |_ tj |¬¦«|_ |€i|_nd|_d|_d|_d|_d|_dS)NrF)Úmaxsize)Ú_connÚ_idÚ_subjectÚ_queueÚ _max_msgsÚ _receivedÚ_cbÚ_futureÚ_closedÚ_pending_msgs_limitÚ_pending_bytes_limitÚasyncioÚQueueÚ_pending_queueÚ_pending_next_msgs_callsÚ _pending_sizeÚ_wait_for_msgs_taskÚ_message_iteratorÚ_jsi) ÚselfÚconnrrrrrrrrs úL/opt/imunify360/venv/lib64/python3.11/site-packages/nats/aio/subscription.pyÚ__init__zSubscription.__init__?sª€ðˆŒ ؈ŒØˆŒ ؈Œ Ø!ˆŒØˆŒØˆŒØˆŒ ؈Œ ð$6ˆÔ Ø$7ˆÔ!Ý29´-ÐHZÐ2[Ñ2[Ô2[ˆÔð ˆ:Ø,.ˆDÔ )Ð )à,0ˆDÔ )ØˆÔØ#'ˆÔ Ø!%ˆÔð6:ˆŒ ˆ ˆ ócó—|jS)z< Returns the subject of the `Subscription`. )r#©r4s r6rzSubscription.subjectfs €ð Œ}Ðr8có—|jS)zX Returns the queue name of the `Subscription` if part of a queue group. )r$r:s r6rzSubscription.queuems €ð Œ{Ðr8úAsyncIterator[Msg]cóF—|jstjd¦«‚|jS)a¨ Retrieves an async iterator for the messages from the subscription. This is only available if a callback isn't provided when creating a subscription. :: nc = await nats.connect() sub = await nc.subscribe('foo') # Use `async for` which implicitly awaits messages async for msg in sub.messages: print('Received', msg) zCcannot iterate over messages with a non iteration subscription type)r2r ÚErrorr:s r6ÚmessageszSubscription.messagests*€ð Ô%ð fÝ”,ÐdÑeÔeÐ eàÔ%Ð%r8có4—|j ¦«S)zw Number of delivered messages by the NATS Server that are being buffered in the pending queue. )r.Úqsizer:s r6Ú pending_msgszSubscription.pending_msgs‰s€ð Ô"×(Ò(Ñ*Ô*Ð*r8có—|jS)zk Size of data sent by the NATS Server that is being buffered in the pending queue. )r0r:s r6Ú pending_byteszSubscription.pending_bytes‘s €ð Ô!Ð!r8có—|jS)zK Number of delivered messages to this subscription so far. )r&r:s r6Ú deliveredzSubscription.delivered™s €ð Œ~Ðr8çð?ÚtimeoutúOptional[float]r cƒó¾‡‡K—dˆˆfd„ }‰jjr tj‚‰jrtjd¦«‚t t¦«¦«} tj |¦«¦«}|‰j |<|ƒd{V—†}‰xj t|j ¦«zc_ ‰j ¦«|‰j  |d¦«S#tj$r%‰jjr tj‚tj‚tj$r‰jjr tj‚‚wxYw#‰j  |d¦«wxYw)a– :params timeout: Time in seconds to wait for next message before timing out. :raises nats.errors.TimeoutError: next_msg can be used to retrieve the next message from a stream of messages using await syntax, this only works when not passing a callback on `subscribe`:: sub = await nc.subscribe('hello') msg = await sub.next_msg(timeout=1) rr c“ól•K—tj‰j ¦«‰¦«ƒd{V—†S©N)r,Úwait_forr.Úget)r4rHs€€r6Ú timed_getz(Subscription.next_msg..timed_get­s;øèè€Ý Ô)¨$Ô*=×*AÒ*AÑ*CÔ*CÀWÑMÔMÐMÐMÐMÐMÐMÐMÐ Mr8z4nats: next_msg cannot be used in async subscriptionsN©rr )r!Ú is_closedr ÚConnectionClosedErrorr'r>rr r,Ú create_taskr/r0ÚlenÚdatar.Ú task_doneÚpopÚ TimeoutErrorÚCancelledError)r4rHrOÚ task_namerÚmsgs`` r6Únext_msgzSubscription.next_msg s‘øøèè€ð Nð Nð Nð Nð Nð Nð Nð Œ:Ô ð /ÝÔ.Ð .à Œ8ð WÝ”,ÐUÑVÔVÐ Vå�™œ‘L”Lˆ ð ?ÝÔ(¨¨©¬Ñ5Ô5ˆFØ7=ˆDÔ )¨)Ñ 4Ø�,�,�,�,�,�,ˆCð Ð Ô ¥# c¤h¡-¤-Ñ /Ð Ô ð Ô × )Ò )Ñ +Ô +Ð +Øà Ô )× -Ò -¨i¸Ñ >Ô >Ð >Ð >øõ!Ô#ð &ð &ð &ØŒzÔ#ð 3ÝÔ2Ð2ÝÔ%Ð %ÝÔ%ð ð ð ØŒzÔ#ð 3ÝÔ2Ð2Ø ð øøøøð Ô )× -Ò -¨i¸Ñ >Ô >Ð >Ð >øøøsÁ.C! Úget_running_looprSÚ_wait_for_msgsr1r(Ú_SubscriptionMessageIteratorr2)r4Úerror_cbs r6Ú_startzSubscription._startÍs¿€ð Œ8ð HÝÔ.¨t¬xÑ8Ô8ð Qݘœ &Ñ)Ô)ð QÝ.5Ô.IÈ$Ì(Ì-Ñ.XÔ.Xð Qõ”lÐ#OÑPÔPÐPå'.Ô'?Ñ'AÔ'A×'MÒ'MÈd×NaÒNaÐbjÑNkÔNkÑ'lÔ'lˆDÔ $Ð $Ð $à Œ\ð Hà ˆDå%AÀ$Ñ%GÔ%GˆDÔ "Ð "Ð "r8cƒóÄK—|jjr tj‚|jjr tj‚|jr tj‚| ¦«ƒd{V—†dS)zU Removes interest in a subject, but will process remaining messages. N) r!rQr rRÚ is_drainingÚConnectionDrainingErrorr)ÚBadSubscriptionErrorÚ_drainr:s r6ÚdrainzSubscription.drainßsmèè€ð Œ:Ô ð /ÝÔ.Ð .Ø Œ:Ô !ð 1ÝÔ0Ð 0Ø Œ<ð .ÝÔ-Ð -Ø�kŠk‰mŒmÐÐÐÐÐÐÐÐÐr8cƒó˜K— |j |j¦«ƒd{V—†|j ¦«ƒd{V—†|jr|j ¦«ƒd{V—†| ¦«|j |j¦«n#tj $r‚wxYw d|_ dS#d|_ wxYw©NT) r!Ú_send_unsubscriber"Úflushr.ÚjoinÚ_stop_processingÚ _remove_subr,rYr)r:s r6rjzSubscription._drainës èè€ð ð”*×.Ò.¨t¬xÑ8Ô8Ð 8Ð 8Ð 8Ð 8Ð 8Ð 8Ð 8ð”*×"Ò"Ñ$Ô$Ð $Ð $Ð $Ð $Ð $Ð $Ð $àÔ"ð 1ðÔ)×.Ò.Ñ0Ô0Ð0Ð0Ð0Ð0Ð0Ð0Ð0ð × !Ò !Ñ #Ô #Ð #ð ŒJ× "Ò " 4¤8Ñ ,Ô ,Ð ,Ð ,øÝÔ%ð ð ð Ø ð øøøð -ð ˆDŒLˆLˆLø˜4ˆDŒLÐ Ð Ð Ð s„BB"Â!CÂ"B3Â3Cà C ÚlimitcƒóÐK—|jjr tj‚|jjr tj‚|jr tj‚||_|dks$|j |krS|j   ¦«r:d|_|  ¦«|j  |j¦«|jjs)|j |j|¬¦«ƒd{V—†dSdS)aX :param limit: Max number of messages to receive before unsubscribing. Removes interest in a subject, remaining messages will be discarded. If `limit` is greater than zero, interest is not immediately removed, rather, interest will be automatically removed after `limit` messages are received. rT)rsN)r!rQr rRrgrhr)rir%r&r.Úemptyrqrrr"Úis_reconnectingrn)r4rss r6Ú unsubscribezSubscription.unsubscribes÷èè€ð Œ:Ô ð /ÝÔ.Ð .Ø Œ:Ô !ð 1ÝÔ0Ð 0Ø Œ<ð .ÝÔ-Ð -àˆŒØ �AŠ:ˆ:˜$œ.¨EÒ1Ð1°dÔ6I×6OÒ6OÑ6QÔ6QÐ1؈DŒLØ × !Ò !Ñ #Ô #Ð #Ø ŒJ× "Ò " 4¤8Ñ ,Ô ,Ð ,àŒzÔ)ð FØ”*×.Ò.¨t¬x¸uÐ.ÑEÔEÐ EÐ EÐ EÐ EÐ EÐ EÐ EÐ EÐ Eð Fð Fr8có¼—|jr2|j ¦«s|j ¦«|jr|j ¦«dSdS)zF Stops the subscription from processing new messages. N)r1ÚdoneÚcancelr2Ú_cancelr:s r6rqzSubscription._stop_processingsj€ð Ô #ð .¨DÔ,D×,IÒ,IÑ,KÔ,Kð .Ø Ô $× +Ò +Ñ -Ô -Ð -Ø Ô !ð -Ø Ô "× *Ò *Ñ ,Ô ,Ð ,Ð ,Ð ,ð -ð -r8cƒó¨K—|js Jd¦«‚ |j ¦«ƒd{V—†}|xjt |j¦«zc_ | |¦«ƒd{V—†nT#t j$rY|j ¦«dSt$r}|r||¦«ƒd{V—†Yd}~nd}~wwxYw|j ¦«n#|j ¦«wxYw|j dkr0|j |j kr |jj r|  ¦«n#t j$rYdSwxYw�Œ?)zz A coroutine to read and process messages if a callback is provided. Should be called as a task. z-_wait_for_msgs can be called only from _startTNr)r'r.rNr0rTrUr,rYrVÚ Exceptionr%r&rurq)r4rdr[Úes r6rbzSubscription._wait_for_msgs'sµèè€ð ŒxÐHÐHÐHÑHÔHˆxð ð Ø Ô/×3Ò3Ñ5Ô5Ð5Ð5Ð5Ð5Ð5Ð5�ØÐ"Ô"¥c¨#¬(¡m¤mÑ3Ð"Ô"ð4àŸ(š( 3™-œ-Ð'Ð'Ð'Ð'Ð'Ð'Ð'Ð'øÝÔ-ðððððÔ'×1Ò1Ñ3Ô3Ð3Ð3Ð3õ!ð*ð*ð*ð ð*Ø&˜h q™kœkÐ)Ð)Ð)Ð)Ð)Ð)Ð)øøøøøøøøøð *øøøðÔ'×1Ò1Ñ3Ô3Ð3Ð3ø�DÔ'×1Ò1Ñ3Ô3Ð3Ð3øøøð”> AÒ%Ð%¨$¬.¸D¼NÒ*JÐ*JÈtÔObÔOhÐ*JØ×)Ò)Ñ+Ô+Ð+øøÝÔ)ð ð ð Ø��ð øøøñ1 s`–AD<ÁA4Á3C"Á4CÂC"ÂD< CÂ(CÂ;C"ÃCÃC"ÃD<Ã"C=Ã=>D<Ä<EÅE)rrrrrrrrrrrrrrrrrr)rr)rr<)rr)rG)rHrIrr ©rr)r)rsr)Ú__name__Ú __module__Ú __qualname__Ú__doc__ÚDEFAULT_SUB_PENDING_MSGS_LIMITÚDEFAULT_SUB_PENDING_BYTES_LIMITr7Úpropertyrrr?rBrDrFr\rerkrjrwrqrb©r8r6rr(s§€€€€€ððð2ØØØ9=Ø+/ØØ"@Ø#Bð%:ð%:ð%:ð%:ð%:ðNðððñ„Xðð ðððñ„Xðð ð&ð&ð&ñ„Xð&ð(ð+ð+ð+ñ„Xð+ðð"ð"ð"ñ„Xð"ððððñ„Xðð +?ð+?ð+?ð+?ð+?ðZHðHðHð$ ð ð ð ð ð ð ð2FðFðFðFðFð4-ð-ð-ð-ð ð ð ð ð ð r8rcó.—eZdZd d„Zd d„Zd d„Zdd „Zd S)rcÚsubrrrcó\—||_|j|_tj¦«|_dSrL)Ú_subr.r$r,ÚFutureÚ_unsubscribed_future)r4r‰s r6r7z%_SubscriptionMessageIterator.__init__Ks)€Ø"%ˆŒ Ø*-Ô*<ˆŒ Ý:A¼.Ñ:JÔ:JˆÔ!Ð!Ð!r8cóp—|j ¦«s|j d¦«dSdSrm)r�ryÚ set_resultr:s r6r{z$_SubscriptionMessageIterator._cancelPs@€ØÔ(×-Ò-Ñ/Ô/ð 7Ø Ô %× 0Ò 0°Ñ 6Ô 6Ð 6Ð 6Ð 6ð 7ð 7r8có—|SrLr‡r:s r6Ú __aiter__z&_SubscriptionMessageIterator.__aiter__Ts€Øˆ r8r cƒólK—tj¦« |j ¦«¦«}||jg}tj|tj¬¦«ƒd{V—†\}}|j}||vr…|j  ¦«|  ¦«}|jxj t|j ¦«zc_ |jdkr$|j|jkr| ¦«|S|j ¦«r| ¦«t&‚)N)Ú return_whenr)r,rarSr$rNr�ÚwaitÚFIRST_COMPLETEDr‹rVÚresultr0rTrUr%r&r{ryrzÚStopAsyncIteration)r4Úget_taskÚtasksÚfinishedÚ_r‰r[s r6Ú __anext__z&_SubscriptionMessageIterator.__anext__Wsèè€ÝÔ+Ñ-Ô-×9Ò9¸$¼+¿/º/Ñ:KÔ:KÑLÔLˆØ'/°Ô1JÐ&KˆÝ#œL¨½GÔr£s5ðð#Ð"Ð"Ð"Ð"Ð"à€€€ðððððððððððððððððÐÐÐÐÐàÐÐÐÐÐðÐÐÐÐÐàð)Ø(Ð(Ð(Ð(Ð(Ð(à!+ÐØ"3Ðð_ð_ð_ð_ð_ñ_ô_ð_ðD !ð!ð!ð!ð!ñ!ô!ð!ð!ð!r8