A helper class for managing CRTTransferFuture
| 633 | |
| 634 | |
| 635 | class CRTTransferCoordinator: |
| 636 | """A helper class for managing CRTTransferFuture""" |
| 637 | |
| 638 | def __init__( |
| 639 | self, transfer_id=None, s3_request=None, exception_translator=None |
| 640 | ): |
| 641 | self.transfer_id = transfer_id |
| 642 | self._exception_translator = exception_translator |
| 643 | self._s3_request = s3_request |
| 644 | self._lock = threading.Lock() |
| 645 | self._exception = None |
| 646 | self._crt_future = None |
| 647 | self._done_event = threading.Event() |
| 648 | |
| 649 | @property |
| 650 | def s3_request(self): |
| 651 | return self._s3_request |
| 652 | |
| 653 | def set_done_callbacks_complete(self): |
| 654 | self._done_event.set() |
| 655 | |
| 656 | def wait_until_on_done_callbacks_complete(self, timeout=None): |
| 657 | self._done_event.wait(timeout) |
| 658 | |
| 659 | def set_exception(self, exception, override=False): |
| 660 | with self._lock: |
| 661 | if not self.done() or override: |
| 662 | self._exception = exception |
| 663 | |
| 664 | def cancel(self): |
| 665 | if self._s3_request: |
| 666 | self._s3_request.cancel() |
| 667 | |
| 668 | def result(self, timeout=None): |
| 669 | if self._exception: |
| 670 | raise self._exception |
| 671 | try: |
| 672 | self._crt_future.result(timeout) |
| 673 | except KeyboardInterrupt: |
| 674 | self.cancel() |
| 675 | self._crt_future.result(timeout) |
| 676 | raise |
| 677 | except Exception as e: |
| 678 | self.handle_exception(e) |
| 679 | finally: |
| 680 | if self._s3_request: |
| 681 | self._s3_request = None |
| 682 | |
| 683 | def handle_exception(self, exc): |
| 684 | translated_exc = None |
| 685 | if self._exception_translator: |
| 686 | try: |
| 687 | translated_exc = self._exception_translator(exc) |
| 688 | except Exception as e: |
| 689 | # Bail out if we hit an issue translating |
| 690 | # and raise the original error. |
| 691 | logger.debug("Unable to translate exception.", exc_info=e) |
| 692 | pass |