-
Notifications
You must be signed in to change notification settings - Fork 108
'await connection.close()' returns once connection thread has also forwarded _STOP_RUNNING_SENTINEL #305
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
'await connection.close()' returns once connection thread has also forwarded _STOP_RUNNING_SENTINEL #305
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -62,10 +62,15 @@ def __init__( | |
| DeprecationWarning, | ||
| ) | ||
|
|
||
| def _stop_running(self): | ||
| def _stop_running(self) -> asyncio.Future: | ||
| self._running = False | ||
| # PEP 661 is not accepted yet, so we cannot type a sentinel | ||
| self._tx.put_nowait(_STOP_RUNNING_SENTINEL) # type: ignore[arg-type] | ||
|
|
||
| function = partial(lambda: _STOP_RUNNING_SENTINEL) | ||
| future = asyncio.get_event_loop().create_future() | ||
|
|
||
| self._tx.put_nowait((future, function)) | ||
|
|
||
| return future | ||
|
|
||
| @property | ||
| def _conn(self) -> sqlite3.Connection: | ||
|
|
@@ -95,16 +100,17 @@ def run(self) -> None: | |
| # futures) | ||
|
|
||
| tx_item = self._tx.get() | ||
| if tx_item is _STOP_RUNNING_SENTINEL: | ||
| break | ||
|
|
||
| future, function = tx_item | ||
|
|
||
| try: | ||
| LOG.debug("executing %s", function) | ||
| result = function() | ||
| LOG.debug("operation %s completed", function) | ||
| future.get_loop().call_soon_threadsafe(set_result, future, result) | ||
|
|
||
| if result is _STOP_RUNNING_SENTINEL: | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This likely caused #369 It might of course be some bad usages of aiosqlite in airflow or in airflow tests, but I think at the very least some diagnostics of the situation and detecting bad use would be good, because we are now totally in the dark what could be the reason - simply pytest session hangs forever with all tests passing. |
||
| break | ||
|
|
||
| except BaseException as e: # noqa B036 | ||
| LOG.debug("returning exception %s", e) | ||
| future.get_loop().call_soon_threadsafe(set_exception, future, e) | ||
|
|
@@ -129,7 +135,7 @@ async def _connect(self) -> "Connection": | |
| self._tx.put_nowait((future, self._connector)) | ||
| self._connection = await future | ||
| except BaseException: | ||
| self._stop_running() | ||
| await self._stop_running() | ||
| self._connection = None | ||
| raise | ||
|
|
||
|
|
@@ -170,7 +176,7 @@ async def close(self) -> None: | |
| LOG.info("exception occurred while closing connection") | ||
| raise | ||
| finally: | ||
| self._stop_running() | ||
| await self._stop_running() | ||
| self._connection = None | ||
|
|
||
| @contextmanager | ||
|
|
||
Uh oh!
There was an error while loading. Please reload this page.