U vi @s\ddlmZmZmZddlmZddlmZmZddl m Z edZ Gddde e Z d S) )ListOptionalTypeVar)abc)CompositeDisposable Disposable) Observable_TcseZdZdZejeejedfdd Zdej ee ej ej dddZ de ej e ej d d d Zdeeed ddZZS)ConnectableObservablezDRepresents an observable that can be connected and disconnected.)sourcesubjectcs&||_d|_d|_||_tdSNF)r has_subscription subscriptionr super__init__)selfr r  __class__[/opt/alt/python38/lib/python3.8/site-packages/reactivex/observable/connectableobservable.pyrs zConnectableObservable.__init__Nobserver schedulerreturncCs|jj||dS)Nr)r subscribe)rrrrrr_subscribe_coresz%ConnectableObservable._subscribe_core)rrcsFjs@d_ddfdd }jjj|d}t|t|_jS)zConnects the observable.TNrcs d_dSr)rrrrrdispose&sz.ConnectableObservable.connect..disposer)rr rr rrr)rrr!rrr rconnects zConnectableObservable.connectr)subscriber_countrcshdgdg|dgdkr2d<dd<dtjtttjtjdfdd }t|S) a]Returns an observable sequence that stays connected to the source indefinitely to the observable sequence. Providing a subscriber_count will cause it to connect() after that many subscriptions occur. A subscriber_count of 0 will result in emissions firing immediately without waiting for subscribers. NrFTrcshdd7<dko$d }||rJ|d<dd<ddfdd }t|S)NrrTrcs$dd8<dd<dS)NrrF)r!r)count is_connectedrrrr!KszFConnectableObservable.auto_connect..subscribe..dispose)rr"r)rrZshould_connectr!Zconnectable_subscriptionr$r%r r#)rrr@s z5ConnectableObservable.auto_connect..subscribe)N)r"r ObserverBaser r SchedulerBaseDisposableBaser )rr#rrr&r auto_connect.s  z"ConnectableObservable.auto_connect)N)N)r)__name__ __module__ __qualname____doc__rZObservableBaser Z SubjectBaserr'rr(r)rr"intr r* __classcell__rrrrr s   r N) typingrrrZ reactivexrZreactivex.disposablerrZ observabler r r rrrrs