Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions DEV.md
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,8 @@ The adapter uses optional capabilities when present:

The parent shell runs on a persistent asyncio loop in the Python main thread. Child shells run on supervised OS threads with their own persistent loops created by the same factory. Kernmini's multi-thread Tokio runtime independently drives transport, queues, output, control, and interrupt futures, so synchronous Python cannot block the engine.

A `SystemExit` raised inside a task leaves the loop rather than the task, so `kernmini._bridge.run_loop` drives every loop and re-enters it: user code cannot end the kernel.

`pyo3-async-runtimes` bridges Python awaitables onto their owning loop. Interrupts cancel async cells through that loop and inject `KeyboardInterrupt` into synchronous Python. A child blocked indefinitely in arbitrary C code cannot be interrupted safely; kernmini does not pretend otherwise.

## Execution and concurrency
Expand Down
5 changes: 4 additions & 1 deletion kernmini/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

import asyncio

from ._bridge import run_loop
from .concur import sidecar, subshell
from .kernelspec import install_kernelspec, install_kernelspec_dir

Expand All @@ -17,7 +18,9 @@ def run_kernel(connection_file, shell_factory, *, loop_factory=None, own_process
from ._native import run_kernel as run_native
if loop_factory is None: loop_factory = _default_loop_factory()
async def run(): await run_native(connection_file, shell_factory, loop_factory, own_process_group=own_process_group)
with asyncio.Runner(loop_factory=loop_factory) as runner: return runner.run(run())
with asyncio.Runner(loop_factory=loop_factory) as runner:
loop = runner.get_loop()
return run_loop(loop, loop.create_task(run()))


def __getattr__(name):
Expand Down
8 changes: 8 additions & 0 deletions kernmini/_bridge.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,14 @@
_current = contextvars.ContextVar("kernmini.execution", default=None)


def run_loop(loop, fut=None):
"Run `loop` until `fut` is done (forever if None), re-entering after a `SystemExit` escapes a task, which asyncio raises out of the loop"
while True:
try: return loop.run_until_complete(fut) if fut is not None else loop.run_forever()
except SystemExit:
if fut is not None and fut.done(): raise


class _IOPub:
def send(self, msg_type, parent=None, content=None, metadata=None, ident=None, buffers=None, **kwargs):
sink = _current.get()
Expand Down
2 changes: 1 addition & 1 deletion src/python.rs
Original file line number Diff line number Diff line change
Expand Up @@ -294,7 +294,7 @@ impl Language for PyLanguage {
let session = PyLanguageSession::with_locals(py, target, locals, Some(thread_loop.clone()))?;
if created.take().unwrap().send(Ok(session)).is_err() { return Ok(()); }
drop(thread_loop);
event_loop.call_method0(py, "run_forever")?;
py.import("kernmini._bridge")?.call_method1("run_loop", (&event_loop,))?;
Ok(())
});
if let Err(error) = outcome
Expand Down
22 changes: 9 additions & 13 deletions tests/ipython_kernel.py
Original file line number Diff line number Diff line change
@@ -1,20 +1,16 @@
"The real ipymini language adapter hosted directly by kernmini."
"The real ipymini language adapter hosted by kernmini's public runner."

import asyncio, sys

from ipymini.shell import MiniShell
from kernmini._native import run_kernel
from kernmini import run_kernel

user_ns, first = {}, True
def shell_factory():
global first
shell = MiniShell(request_input=lambda *_: "", user_ns=user_ns, use_singleton=first)
first = False
return shell

async def main():
user_ns, first = {}, True
def shell_factory():
nonlocal first
shell = MiniShell(request_input=lambda *_: "", user_ns=user_ns, use_singleton=first)
first = False
return shell
await run_kernel(sys.argv[-1], shell_factory, asyncio.new_event_loop)


if __name__ == "__main__":
with asyncio.Runner() as runner: runner.run(main())
if __name__ == "__main__": run_kernel(sys.argv[-1], shell_factory, loop_factory=asyncio.new_event_loop)
6 changes: 6 additions & 0 deletions tests/test_rust_ipython_kernel.py
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,12 @@ async def test_ipython_story():
background = await qs.jmsg_for('stream', pred=lambda m: parent_id(m) == task_id, queue='iopub', timeout=5)
assert background['content']['text'] == 'background\n'

bail = "async def bail(): raise SystemExit('bail')\nres = await asyncio.gather(bail(), return_exceptions=True)\ntype(res[0]).__name__"
for kw in ({}, dict(subshell_id='sidecar')):
msgs = await _run(kc, bail, **kw)
assert _one(msgs, 'execute_reply')['content']['status'] == 'ok'
assert _one(msgs, 'execute_result')['content']['data']['text/plain'] == "'SystemExit'"

sleeper = kc.run("print('sleeping', flush=True)\nawait asyncio.sleep(.3)", timeout=5)
sleeper_msgs = await _until_stream(sleeper, 'sleeping\n')
complete = await kc.cmd.complete(code='x.rea', cursor_pos=5, timeout=5)
Expand Down