Running an async function in an worker thread/loop

Viewed 233

I'm trying to write code that enables using asyncpg from mostly sync code (to avoid duplication). For some very strange reason, the coroutine Database.test() will execute and return in my worker eventloop/thread. The future seems to work correctly. But connecting to a database with asyncpg will just hang. Any clue as to why?

Also, maybe I should use asyncio.run() instead.

from threading import Thread
import asyncio
import asyncpg

class AsyncioWorkerThread(Thread):

    def __init__(self, *args, daemon=True, loop=None, **kwargs):
        super().__init__(*args, daemon=daemon, **kwargs)
        self.loop = loop or asyncio.new_event_loop()
        self.running = False

    def run(self):
        self.running = True
        self.loop.run_forever()

    def submit(self, coro):
        fut = asyncio.run_coroutine_threadsafe(coro, loop=self.loop)
        return fut.result()

    def stop(self):
        self.loop.call_soon_threadsafe(self.loop.stop)
        self.join()
        self.running = False

class Database:
    async def test(self):
        print('In test')
        await asyncio.sleep(5)

    async def connect(self):
        # Put in your db credentials here
        # pg_user = ''
        # pg _password = ''
        # pg_host = ""
        # pg_port = 20
        # pg_db
        connection_uri = f'postgres://{pg_user}:{pg_password}@{pg_host}:{pg_port}/{pg_db}'
        self.connection_pool = await asyncpg.create_pool(
            connection_uri, min_size=5, max_size=10)

if __name__ == "__main__":
    db = Database()
    worker = AsyncioWorkerThread()
    worker.start()
    worker.submit(db.test())  # Works future returns correctly
    worker.submit(db.connect())  # Hangs, thread never manages to acquire

0 Answers
Related