@@ -2,18 +2,18 @@ use std::future::IntoFuture;
22use std:: sync:: Arc ;
33use std:: time:: Duration ;
44
5- use futures:: { stream , StreamExt , TryStreamExt } ;
5+ use futures:: { StreamExt , TryStreamExt , stream } ;
66use opsqueue:: {
7+ E ,
78 common:: errors:: {
8- IncorrectUsage , LimitIsZero ,
99 E :: { self , L , R } ,
10+ IncorrectUsage , LimitIsZero ,
1011 } ,
1112 consumer:: client:: InternalConsumerClientError ,
1213 object_store:: {
1314 ChunkRetrievalError , ChunkStorageError , ChunkType , NewObjectStoreClientError ,
1415 ObjectStoreClient ,
1516 } ,
16- E ,
1717} ;
1818use pyo3:: {
1919 create_exception,
@@ -350,18 +350,23 @@ impl ConsumerClient {
350350 let submission_prefix = chunk. submission_prefix . clone ( ) ;
351351 let chunk_index = chunk. chunk_index ;
352352 tracing:: debug!(
353- "Running fun for chunk: submission_id={:?}, chunk_index={:?}, submission_prefix={:?}" ,
354- submission_id,
355- chunk_index,
356- & submission_prefix
357- ) ;
358- let res = Python :: attach ( |py| {
353+ "Running fun for chunk: submission_id={:?}, chunk_index={:?}, submission_prefix={:?}" ,
354+ submission_id,
355+ chunk_index,
356+ & submission_prefix
357+ ) ;
358+ let res = Python :: with_gil ( |py| {
359359 let res = unbound_fun. bind ( py) . call1 ( ( chunk, ) ) ?;
360360 res. extract ( )
361361 } ) ;
362362 match res {
363363 Ok ( res) => {
364- tracing:: debug!( "Completing chunk: submission_id={:?}, chunk_index={:?}, submission_prefix={:?}" , submission_id, chunk_index, & submission_prefix) ;
364+ tracing:: debug!(
365+ "Completing chunk: submission_id={:?}, chunk_index={:?}, submission_prefix={:?}" ,
366+ submission_id,
367+ chunk_index,
368+ & submission_prefix
369+ ) ;
365370 self . complete_chunk_gilless (
366371 submission_id,
367372 submission_prefix. clone ( ) ,
@@ -374,15 +379,20 @@ impl ConsumerClient {
374379 CError ( R ( R ( e) ) ) => CError ( R ( R ( R ( L ( e) ) ) ) ) ,
375380 } ) ?;
376381 tracing:: debug!(
377- "Completed chunk: submission_id={:?}, chunk_index={:?}, submission_prefix={:?}" ,
378- submission_id,
379- chunk_index,
380- & submission_prefix
381- ) ;
382+ "Completed chunk: submission_id={:?}, chunk_index={:?}, submission_prefix={:?}" ,
383+ submission_id,
384+ chunk_index,
385+ & submission_prefix
386+ ) ;
382387 }
383388 Err ( failure) => {
384389 let failure_str = crate :: common:: format_pyerr ( & failure) ;
385- tracing:: warn!( "Failing chunk: submission_id={:?}, chunk_index={:?}, submission_prefix={:?}, reason: {failure_str}" , submission_id, chunk_index, & submission_prefix) ;
390+ tracing:: warn!(
391+ "Failing chunk: submission_id={:?}, chunk_index={:?}, submission_prefix={:?}, reason: {failure_str}" ,
392+ submission_id,
393+ chunk_index,
394+ & submission_prefix
395+ ) ;
386396 self . fail_chunk_gilless (
387397 submission_id,
388398 submission_prefix. clone ( ) ,
@@ -394,11 +404,11 @@ impl ConsumerClient {
394404 CError ( R ( e) ) => CError ( R ( R ( R ( L ( e) ) ) ) ) ,
395405 } ) ?;
396406 tracing:: warn!(
397- "Failed chunk: submission_id={:?}, chunk_index={:?}, submission_prefix={:?}" ,
398- submission_id,
399- chunk_index,
400- & submission_prefix
401- ) ;
407+ "Failed chunk: submission_id={:?}, chunk_index={:?}, submission_prefix={:?}" ,
408+ submission_id,
409+ chunk_index,
410+ & submission_prefix
411+ ) ;
402412
403413 // On exceptions that are not PyExceptions (but PyBaseExceptions), like KeyboardInterrupt etc, return.
404414 if !Python :: attach ( |py| failure. is_instance_of :: < PyException > ( py) ) {
0 commit comments