dask / dask/distributed

handle `distributed.core.Server` startup and shutdown excellently

Open
#6,616 1 comment 1 reaction 0 assignees View on GitHub
asyncio deprecation enhancement stability
Dominant language
Python
Stars
1.7k
Forks
778
Avg merge
2h 50m
Merged PRs (30d)
3

Description

Server startup and shutdown is currently confusing and error prone see https://github.com/dask/distributed/pull/6615 and [`a8244bd` (#6603)](https://github.com/dask/distributed/pull/6603/commits/a8244bde99182ffb7abd6508285d566414f95b22)

patterns evolving concurrent and re-entrant close and cancellation are prone to deadlocks:

```python
try:
self.comm.read() # close call cancels this task and waits for this task to finish
finally:
await self.close() # this waits for close to finish
```
eg https://github.com/dask/distributed/blob/bc04d0e29077c675f24b867e435a6aef3f0652cd/distributed/worker.py#L1201-L1210

I think a pattern where only the task that calls `async with Server(...): ...` are allowed to call `await self.finish()` or `await self.close()`

a sketch here

```python
class Server:
def __init__(self):
self.__close_done = asyncio.Event()
self.__start_event = asyncio.Event()
self.__close_event = asyncio.Event()

def request_close(self):
self.__start_event.set()
self.__close_event.set()

async def __lifecycle(self):
try:
await self.__start_event.wait()
async with self.listen(), self.open_rpc_pool():
try:
await self.start()
self.__start_event.set()
await self.__close_event.wait()
await self.close()
finally:
v = self.abort() # abort comms by calling socket.close()
assert v is None # abort must not be an async def
finally:
self.__close_done.set()

async def __aenter__(self):
self.__parent_task = asyncio.current_task()
self.__lifecyle_task = asyncio.create_task(self.__lifecyle())
await self.__start_event.wait()

def __await__(self):
warnings.warn("await Server() is deprecated, use async with Server()")
# ??? some background task magic here

async def __aexit__(self):
self.request_close()
await self.finished()

async def close(self):
try:
assert asyncio.current_task() is self.__lifecyle_task
# close comms, wait for tasks to cancel
finally:
v = self.abort() # abort comms by calling socket.close()
assert v is None

async def finished(self):
assert asyncio.current_task() is self.__parent_task
await self.__close_done.wait()
```

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.