U vâ’ižã@sddlmZmZmZmZmZddlmZmZddlm Z edƒZ eeeee geee fdœdd„Z egeefeee geee fdœd d „Z eeeegeefeee geee fd œd d „Zdeeeeee geee fdœdd„Zddd d gZdS)é)ÚAnyÚCallableÚListÚOptionalÚTypeVar)Ú ObservableÚcompose)Ú operatorsÚ_T)Ú boundariesÚreturncCstt |¡t t ¡¡ƒS©N)rÚopsZwindowÚflat_mapÚto_list)r ©rúL/opt/alt/python38/lib/python3.8/site-packages/reactivex/operators/_buffer.pyÚbuffer_ s þr)Úclosing_mapperr cCstt |¡t t ¡¡ƒSr )rrZ window_whenrr)rrrrÚ buffer_when_s þr)Úopeningsrr cCstt ||¡t t ¡¡ƒSr )rrZ window_togglerr)rrrrrÚbuffer_toggle_s  þrN)ÚcountÚskipr cs&tttttdœ‡‡fdd„ }|S)a9Projects each element of an observable sequence into zero or more buffers which are produced based on element count information. Examples: >>> res = buffer_with_count(10)(xs) >>> res = buffer_with_count(10, 1)(xs) Args: count: Length of each buffer. skip: [Optional] Number of elements to skip between creation of consecutive buffers. If not provided, defaults to the count. Returns: A function that takes an observable source and returns an observable sequence of buffers. )Úsourcer cs^ˆdkr ˆ‰tttttdœdd„}tttdœdd„}| t ˆˆ¡t |¡t |¡¡S)N)Úvaluer cSs| t ¡¡Sr )Úpiperr©rrrrÚmapper?sÿz=buffer_with_count_..buffer_with_count..mappercSs t|ƒdkS)Nr)ÚlenrrrrÚ predicateDsz@buffer_with_count_..buffer_with_count..predicate) rr rÚboolrrZwindow_with_countrÚfilter)rrr ©rrrrÚbuffer_with_count9s ýz-buffer_with_count_..buffer_with_count)rr r)rrr$rr#rÚbuffer_with_count_$s"r%)N)ÚtypingrrrrrZ reactivexrrr rr rrrÚintr%Ú__all__rrrrÚs( þ þ þ ÿþ ,