diff --git a/docs/plugin-api-additions/async.md b/docs/plugin-api-additions/async.md new file mode 100644 index 00000000..8d4160ff --- /dev/null +++ b/docs/plugin-api-additions/async.md @@ -0,0 +1,30 @@ +# Optional Python asyncio integration + +`run_async(coroutine, context=None, interval=10)` returns an AsyncTask with +`done()`, `cancel()` and `result()`. Each invocation owns a separate event loop +pumped by an existing main-loop timer. A ready stop callback ensures selector +polling never blocks the UI. The captured context is restored after every pump; +a closed context cancels the work. No process-wide event-loop policy is changed. + +Existing synchronous scripts are unaffected. Use `asyncio.get_running_loop()` +inside the coroutine. Child tasks use the same loop/context and are cancelled +when the main coroutine finishes or the script unloads. Cancellation cleanup is +bounded; tasks must not suppress cancellation indefinitely. Cleanup uses the +captured context while it remains open; after it closes, cancellation cleanup uses +the current valid context. Use explicit Context methods for context-specific +output in finally blocks, and do not issue context-dependent commands there. This provides normal +asyncio sockets/timers, not Windows subprocess support or a new HTTP library. +Asyncio hostname resolution can use its executor; use a resolver appropriate to +your library if you require no worker threads. Blocking work inside a coroutine +still blocks the client and must be avoided. + +```python +import asyncio +import zoitechat + +async def later(): + await asyncio.sleep(1) + zoitechat.prnt('Still in the context where this task was started') + +task = zoitechat.run_async(later()) +``` diff --git a/plugins/python/_zoitechat_async.py b/plugins/python/_zoitechat_async.py new file mode 100644 index 00000000..be74118f --- /dev/null +++ b/plugins/python/_zoitechat_async.py @@ -0,0 +1,100 @@ +"""Opt-in asyncio tasks pumped without blocking the client main loop.""" +import asyncio +import operator +import traceback +import _zoitechat as api + +__all__ = ['AsyncTask', 'run_async'] + + +class AsyncTask: + """A script-owned task and isolated event loop; cancel() is cooperative.""" + def __init__(self, coroutine, plugin, context, interval): + self._plugin = plugin + self._context = context + self._loop = asyncio.new_event_loop() + self._task = self._loop.create_task(coroutine) + self._closed = False + self._unload_hook = api.hook_unload(lambda data: self._close()) + self._timer_hook = api.hook_timer(interval, self._tick) + + def done(self): + return self._task.done() + + def cancel(self): + if not self._closed: + return self._task.cancel() + return False + + def result(self): + return self._task.result() + + def _pump(self): + # A ready stop callback forces the selector's timeout to zero. + self._loop.call_soon(self._loop.stop) + self._loop.run_forever() + + def _close(self): + if self._closed: + return + self._closed = True + previous = api.get_context() + self._context.set() + try: + pending = asyncio.all_tasks(self._loop) + for task in pending: + task.cancel() + # Bound cleanup work; tasks must cooperate with cancellation. + for unused in range(3): + if not any(not task.done() for task in pending): + break + self._pump() + for task in pending: + if not task.done(): + self._loop.call_exception_handler({ + 'message': 'Script task ignored cancellation during unload', + 'task': task, + }) + finally: + self._loop.close() + previous.set() + + def _tick(self, userdata): + if self._closed: + return False + previous = api.get_context() + if not self._context.set(): + self._close() + self._plugin.remove_hook(self._unload_hook) + return False + try: + self._pump() + if not self._task.done(): + return True + try: + self._task.result() + except asyncio.CancelledError: + pass + except Exception: + traceback.print_exc() + self._close() + self._plugin.remove_hook(self._unload_hook) + return False + except Exception: + traceback.print_exc() + self._close() + self._plugin.remove_hook(self._unload_hook) + return False + finally: + previous.set() + + +def run_async(coroutine, context=None, interval=10): + """Schedule a coroutine in a captured context; return an AsyncTask handle.""" + interval = operator.index(interval) + if interval < 1 or interval > 2147483647: + raise ValueError('interval must be a positive C int in milliseconds') + if not asyncio.iscoroutine(coroutine): + raise TypeError('run_async expects a coroutine object') + plugin = api.__get_current_plugin() + return AsyncTask(coroutine, plugin, context or api.get_context(), interval)