Skip to content

Commit 4de0533

Browse files
authored
fix: coordinate concurrent async runner operations
Submit concurrent work without holding completion locks, reject loop-thread synchronous reentry, and define close-during-startup cancellation.
1 parent 4259699 commit 4de0533

2 files changed

Lines changed: 65 additions & 31 deletions

File tree

‎posthog/_async_utils.py‎

Lines changed: 36 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -57,47 +57,52 @@ def __init__(self) -> None:
5757
self._startup_error: BaseException | None = None
5858
self._close_requested = False
5959
self._lock = threading.Lock()
60-
self._operation_lock = threading.Lock()
6160

6261
def run(self, awaitable: Awaitable[Any]) -> Any:
63-
with self._operation_lock:
62+
if threading.current_thread() is self._thread:
63+
raise RuntimeError("cannot synchronously run from the runner thread")
64+
65+
while True:
6466
loop = self._ensure_loop()
65-
future = asyncio.run_coroutine_threadsafe(
66-
self._await_result(awaitable), loop
67-
)
68-
return future.result()
67+
with self._lock:
68+
if loop is self._loop and not self._close_requested:
69+
future = asyncio.run_coroutine_threadsafe(
70+
self._await_result(awaitable), loop
71+
)
72+
break
73+
return future.result()
6974

7075
def close(self) -> None:
7176
current = threading.current_thread()
7277
with self._lock:
73-
if current is self._thread and self._loop is not None:
74-
self._loop.call_soon(self._loop.stop)
78+
loop = self._loop
79+
thread = self._thread
80+
if thread is None:
7581
return
76-
77-
with self._operation_lock:
78-
with self._lock:
79-
loop = self._loop
80-
thread = self._thread
81-
if thread is None:
82-
return
83-
if loop is None:
84-
self._close_requested = True
85-
else:
86-
self._loop = None
87-
self._thread = None
88-
self._closing_threads.add(thread)
89-
82+
self._close_requested = True
9083
if loop is None:
84+
self._startup_error = RuntimeError("runner closed during startup")
85+
else:
86+
self._loop = None
87+
self._thread = None
88+
self._closing_threads.add(thread)
89+
90+
if loop is None:
91+
if thread is not current:
9192
thread.join()
92-
return
93+
return
9394

94-
if loop.is_closed():
95-
with self._lock:
96-
self._closing_threads.discard(thread)
97-
return
95+
if loop.is_closed():
96+
with self._lock:
97+
self._closing_threads.discard(thread)
98+
return
99+
100+
if thread is current:
101+
loop.call_soon(loop.stop)
102+
return
98103

99-
loop.call_soon_threadsafe(loop.stop)
100-
thread.join()
104+
loop.call_soon_threadsafe(loop.stop)
105+
thread.join()
101106

102107
def owns_thread(self, thread: threading.Thread) -> bool:
103108
with self._lock:
@@ -152,6 +157,8 @@ def _run_loop(self) -> None:
152157
with self._lock:
153158
self._loop = loop
154159
close_requested = self._close_requested
160+
if close_requested and self._startup_error is None:
161+
self._startup_error = RuntimeError("runner closed during startup")
155162
self._started.set()
156163

157164
if close_requested:

‎posthog/test/test_async_utils.py‎

Lines changed: 29 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,9 +32,11 @@ def create_loop():
3232
return _ContextEventLoop()
3333

3434
def run():
35+
awaitable = asyncio.sleep(0)
3536
try:
36-
runner.run(asyncio.sleep(0))
37+
runner.run(awaitable)
3738
except BaseException as error:
39+
awaitable.close()
3840
run_errors.append(error)
3941

4042
with mock.patch(
@@ -52,6 +54,31 @@ def run():
5254

5355
self.assertFalse(run_thread.is_alive())
5456
self.assertFalse(close_thread.is_alive())
55-
self.assertEqual(run_errors, [])
57+
self.assertEqual(len(run_errors), 1)
58+
self.assertRegex(str(run_errors[0]), "closed during startup")
5659
self.assertIsNone(runner._thread)
5760
self.assertIsNone(runner._loop)
61+
62+
def test_run_from_runner_thread_fails_instead_of_deadlocking(self):
63+
runner = _BackgroundEventLoopRunner()
64+
65+
async def reenter():
66+
awaitable = asyncio.sleep(0)
67+
try:
68+
with self.assertRaisesRegex(RuntimeError, "runner thread"):
69+
runner.run(awaitable)
70+
finally:
71+
awaitable.close()
72+
73+
runner.run(reenter())
74+
runner.close()
75+
76+
def test_close_from_runner_thread_allows_fresh_loop(self):
77+
runner = _BackgroundEventLoopRunner()
78+
79+
async def close_runner():
80+
runner.close()
81+
82+
runner.run(close_runner())
83+
runner.run(asyncio.sleep(0))
84+
runner.close()

0 commit comments

Comments
 (0)