4545async def retry_hook (settings , ** kwargs ):
4646 if settings ["hook" ]:
4747 if asyncio .iscoroutine (settings ["hook" ]):
48- await settings ["hook" ](
49- retry_count = settings ["count" ] - 1 ,
50- location_mode = settings ["mode" ],
51- ** kwargs
52- )
48+ await settings ["hook" ](retry_count = settings ["count" ] - 1 , location_mode = settings ["mode" ], ** kwargs )
5349 else :
54- settings ["hook" ](
55- retry_count = settings ["count" ] - 1 ,
56- location_mode = settings ["mode" ],
57- ** kwargs
58- )
50+ settings ["hook" ](retry_count = settings ["count" ] - 1 , location_mode = settings ["mode" ], ** kwargs )
5951
6052
6153async def is_checksum_retry (response ):
@@ -70,9 +62,9 @@ async def is_checksum_retry(response):
7062 await response .http_response .load_body () # Load the body in memory and close the socket
7163 except (StreamClosedError , StreamConsumedError ):
7264 pass
73- computed_md5 = response .http_request .headers .get (
74- "content-md5" , None
75- ) or encode_base64 ( calculate_content_md5 ( response . http_response . body ()))
65+ computed_md5 = response .http_request .headers .get ("content-md5" , None ) or encode_base64 (
66+ calculate_content_md5 ( response . http_response . body ())
67+ )
7668 if response .http_response .headers ["content-md5" ] != computed_md5 :
7769 return True
7870 return False
@@ -118,50 +110,36 @@ async def send(self, request: "PipelineRequest") -> "PipelineResponse":
118110 data_stream_total = request .context .options .pop ("data_stream_total" , None )
119111 download_stream_current = request .context .get ("download_stream_current" )
120112 if download_stream_current is None :
121- download_stream_current = request .context .options .pop (
122- "download_stream_current" , None
123- )
113+ download_stream_current = request .context .options .pop ("download_stream_current" , None )
124114 upload_stream_current = request .context .get ("upload_stream_current" )
125115 if upload_stream_current is None :
126- upload_stream_current = request .context .options .pop (
127- "upload_stream_current" , None
128- )
116+ upload_stream_current = request .context .options .pop ("upload_stream_current" , None )
129117
130- response_callback = request .context .get (
131- "response_callback"
132- ) or request . context . options . pop ( "raw_response_hook" , self . _response_callback )
118+ response_callback = request .context .get ("response_callback" ) or request . context . options . pop (
119+ "raw_response_hook" , self . _response_callback
120+ )
133121
134122 response = await self .next .send (request )
135- will_retry = is_retry (
136- response , request .context .options .get ("mode" )
137- ) or await is_checksum_retry (response )
123+ will_retry = is_retry (response , request .context .options .get ("mode" )) or await is_checksum_retry (response )
138124
139125 # Auth error could come from Bearer challenge, in which case this request will be made again
140126 is_auth_error = response .http_response .status_code == 401
141127 should_update_counts = not (will_retry or is_auth_error )
142128
143129 if should_update_counts and download_stream_current is not None :
144- download_stream_current += int (
145- response .http_response .headers .get ("Content-Length" , 0 )
146- )
130+ download_stream_current += int (response .http_response .headers .get ("Content-Length" , 0 ))
147131 if data_stream_total is None :
148132 content_range = response .http_response .headers .get ("Content-Range" )
149133 if content_range :
150- data_stream_total = int (
151- content_range .split (" " , 1 )[1 ].split ("/" , 1 )[1 ]
152- )
134+ data_stream_total = int (content_range .split (" " , 1 )[1 ].split ("/" , 1 )[1 ])
153135 else :
154136 data_stream_total = download_stream_current
155137 elif should_update_counts and upload_stream_current is not None :
156- upload_stream_current += int (
157- response .http_request .headers .get ("Content-Length" , 0 )
158- )
138+ upload_stream_current += int (response .http_request .headers .get ("Content-Length" , 0 ))
159139 for pipeline_obj in [request , response ]:
160140 if hasattr (pipeline_obj , "context" ):
161141 pipeline_obj .context ["data_stream_total" ] = data_stream_total
162- pipeline_obj .context ["download_stream_current" ] = (
163- download_stream_current
164- )
142+ pipeline_obj .context ["download_stream_current" ] = download_stream_current
165143 pipeline_obj .context ["upload_stream_current" ] = upload_stream_current
166144 if response_callback :
167145 if asyncio .iscoroutine (response_callback ):
@@ -190,9 +168,7 @@ async def send(self, request):
190168 while retries_remaining :
191169 try :
192170 response = await self .next .send (request )
193- if is_retry (
194- response , retry_settings ["mode" ]
195- ) or await is_checksum_retry (response ):
171+ if is_retry (response , retry_settings ["mode" ]) or await is_checksum_retry (response ):
196172 retries_remaining = self .increment (
197173 retry_settings ,
198174 request = request .http_request ,
@@ -211,9 +187,7 @@ async def send(self, request):
211187 except AzureError as err :
212188 if isinstance (err , AzureSigningError ):
213189 raise
214- retries_remaining = self .increment (
215- retry_settings , request = request .http_request , error = err
216- )
190+ retries_remaining = self .increment (retry_settings , request = request .http_request , error = err )
217191 if retries_remaining :
218192 await retry_hook (
219193 retry_settings ,
@@ -275,9 +249,7 @@ def __init__(
275249 self .initial_backoff = initial_backoff
276250 self .increment_base = increment_base
277251 self .random_jitter_range = random_jitter_range
278- super (ExponentialRetry , self ).__init__ (
279- retry_total = retry_total , retry_to_secondary = retry_to_secondary , ** kwargs
280- )
252+ super (ExponentialRetry , self ).__init__ (retry_total = retry_total , retry_to_secondary = retry_to_secondary , ** kwargs )
281253
282254 def get_backoff_time (self , settings : Dict [str , Any ]) -> float :
283255 """
@@ -290,14 +262,8 @@ def get_backoff_time(self, settings: Dict[str, Any]) -> float:
290262 :rtype: int or None
291263 """
292264 random_generator = random .Random ()
293- backoff = self .initial_backoff + (
294- 0 if settings ["count" ] == 0 else pow (self .increment_base , settings ["count" ])
295- )
296- random_range_start = (
297- backoff - self .random_jitter_range
298- if backoff > self .random_jitter_range
299- else 0
300- )
265+ backoff = self .initial_backoff + (0 if settings ["count" ] == 0 else pow (self .increment_base , settings ["count" ]))
266+ random_range_start = backoff - self .random_jitter_range if backoff > self .random_jitter_range else 0
301267 random_range_end = backoff + self .random_jitter_range
302268 return random_generator .uniform (random_range_start , random_range_end )
303269
@@ -335,9 +301,7 @@ def __init__(
335301 """
336302 self .backoff = backoff
337303 self .random_jitter_range = random_jitter_range
338- super (LinearRetry , self ).__init__ (
339- retry_total = retry_total , retry_to_secondary = retry_to_secondary , ** kwargs
340- )
304+ super (LinearRetry , self ).__init__ (retry_total = retry_total , retry_to_secondary = retry_to_secondary , ** kwargs )
341305
342306 def get_backoff_time (self , settings : Dict [str , Any ]) -> float :
343307 """
@@ -352,28 +316,18 @@ def get_backoff_time(self, settings: Dict[str, Any]) -> float:
352316 random_generator = random .Random ()
353317 # the backoff interval normally does not change, however there is the possibility
354318 # that it was modified by accessing the property directly after initializing the object
355- random_range_start = (
356- self .backoff - self .random_jitter_range
357- if self .backoff > self .random_jitter_range
358- else 0
359- )
319+ random_range_start = self .backoff - self .random_jitter_range if self .backoff > self .random_jitter_range else 0
360320 random_range_end = self .backoff + self .random_jitter_range
361321 return random_generator .uniform (random_range_start , random_range_end )
362322
363323
364324class AsyncStorageBearerTokenCredentialPolicy (AsyncBearerTokenCredentialPolicy ):
365325 """Custom Bearer token credential policy for following Storage Bearer challenges"""
366326
367- def __init__ (
368- self , credential : "AsyncTokenCredential" , audience : str , ** kwargs : Any
369- ) -> None :
370- super (AsyncStorageBearerTokenCredentialPolicy , self ).__init__ (
371- credential , audience , ** kwargs
372- )
327+ def __init__ (self , credential : "AsyncTokenCredential" , audience : str , ** kwargs : Any ) -> None :
328+ super (AsyncStorageBearerTokenCredentialPolicy , self ).__init__ (credential , audience , ** kwargs )
373329
374- async def on_challenge (
375- self , request : "PipelineRequest" , response : "PipelineResponse"
376- ) -> bool :
330+ async def on_challenge (self , request : "PipelineRequest" , response : "PipelineResponse" ) -> bool :
377331 try :
378332 auth_header = response .http_response .headers .get ("WWW-Authenticate" )
379333 challenge = StorageHttpChallenge (auth_header )
0 commit comments