Jupyter kernels, owned by the server process, one per notebook session.
Ported from SUNA's python/suna_kernel/bridge.py (github.com/idossha/SUNA,
docs/dev/ARCHITECTURE.md §16.2) — the protocol translation, the nbformat-verbatim
output rule, the request attribution and the fatal error codes are that file's
design. What changed is the topology.
SUNA is a desktop app whose kernel runs on the user's own machine, so the
bridge is a child process the Electron main process speaks to over stdio. In
TI-Toolbox the environment that must be loaded is the one inside the
container — SimNIBS Python with tit importable — and the server already
runs there. So there is no bridge process and no pipe: this module drives
jupyter_client.manager.KernelManager in-process, and the pipe SUNA framed
over stdin/stdout is the WebSocket in routes/kernels.py instead.
The events are SUNA's, unchanged, because they are the shape the ported
renderer already understands:
{"type": "ready", "kernel": {...}}
{"type": "status", "state": "busy"|"idle"|"starting"|"dead"}
{"type": "input", "reqId": "r1", "executionCount": 3}
{"type": "output", "reqId": "r1", "output": {<nbformat output>}}
{"type": "clear", "reqId": "r1", "wait": false}
{"type": "reply", "reqId": "r1", "status": "ok", "executionCount": 3}
{"type": "fatal", "code": "...", "message": "..."}
plus two SUNA does not have, because SUNA's cells are a CodeMirror driven by a
language server and these cells are driven by the kernel itself:
{"type": "complete", "reqId": "c1", "matches": [...],
"cursorStart": 18, "cursorEnd": 24, "metadata": {...}}
{"type": "inspect", "reqId": "i1", "found": true, "text": "..."}
A kernel that has executed the notebook knows what tit. holds better than
any static analyser could -- it is holding the object. That is why completion
is a kernel round trip here rather than an LSP: jedi is already inside
ipykernel, and the namespace it completes against is the live one.
Two limits exist because a kernel here is a container-wide resource rather
than one user's laptop: at most MAX_KERNELS run at once, and a kernel
nobody has touched for IDLE_TIMEOUT_SECONDS is shut down. A FEM run and
four abandoned kernels do not get to share the same RAM.
There is no sandbox. A kernel executes arbitrary user code as the container's
own user, with the project mounted. That is the same trust boundary the job
runners already have, and nothing may call it a sandbox.
KernelError
KernelError(code: str, message: str)
Bases: Exception
A kernel operation failed, with one of the codes the UI switches on.
The codes are SUNA's: no-jupyter-client, no-kernelspec,
start-failed, op-failed — plus too-many-kernels and
no-such-kernel, which only exist because this server is shared and
SUNA's desktop was not.
Source code in tit/server/kernels.py
| def __init__(self, code: str, message: str) -> None:
super().__init__(message)
self.code = code
self.message = message
|
KernelSession
dataclass
KernelSession(id: str, kernel_name: str, cwd: str, manager: Any, client: Any, display_name: str = '', language: str = '', started_at: float = time(), last_used: float = time(), in_flight: int = 0, execution_state: str = 'idle', pending: dict[str, str] = dict(), finishing: dict[str, dict[str, Any]] = dict(), queries: dict[str, tuple[str, str]] = dict(), listeners: list[Listener] = list(), lock: Lock = Lock(), stopping: Event = Event(), pumps: list[Thread] = list())
One live kernel and the bookkeeping the routes need to talk to it.
KernelRegistry
KernelRegistry(*, max_kernels: int = MAX_KERNELS, idle_timeout: float = IDLE_TIMEOUT_SECONDS, clock: Callable[[], float] = time)
Every kernel this server process owns.
One instance lives on the FastAPI app; the routes reach it through
get_kernel_registry. It is deliberately a plain object with a lock
rather than anything async — jupyter_client's channels are threads,
and pretending otherwise would only move the threads somewhere less
visible.
Source code in tit/server/kernels.py
| def __init__(
self,
*,
max_kernels: int = MAX_KERNELS,
idle_timeout: float = IDLE_TIMEOUT_SECONDS,
clock: Callable[[], float] = time.time,
) -> None:
self._sessions: dict[str, KernelSession] = {}
self._lock = threading.Lock()
self.max_kernels = max_kernels
self.idle_timeout = idle_timeout
#: Injectable so a test can age a kernel without sleeping.
self._clock = clock
#: Slots reserved by a start that has not registered its session yet.
#: Counted against the cap: startup takes seconds, and without this N
#: concurrent starts all passed a check none of them had yet answered.
self._starting = 0
|
start
Start a kernel, or raise :class:KernelError saying why not.
Source code in tit/server/kernels.py
| def start(self, *, cwd: str | Path, kernel_name: str = DEFAULT_KERNEL_NAME) -> KernelSession:
"""Start a kernel, or raise :class:`KernelError` saying why not."""
self.reap_idle()
with self._lock:
if len(self._sessions) + self._starting >= self.max_kernels:
raise KernelError(
"too-many-kernels",
f"{self.max_kernels} kernels are already running, which is the limit "
f"for one TI-Toolbox container. Shut one down and try again.",
)
self._starting += 1
try:
return self._start_reserved(cwd, kernel_name)
finally:
# Released whether the start succeeded, raised, or the session is
# already registered: the reservation only has to cover the gap.
with self._lock:
self._starting -= 1
|
shutdown_all
Called when the server stops, so no kernel outlives its owner.
Source code in tit/server/kernels.py
| def shutdown_all(self) -> None:
"""Called when the server stops, so no kernel outlives its owner."""
with self._lock:
sessions = list(self._sessions.values())
self._sessions.clear()
for session in sessions:
self._shutdown_session(session)
|
reap_idle
Shut down kernels nobody has used lately. Returns their ids.
Idle means no request in flight and not executing — a cell that runs
for longer than the timeout is work, not neglect, and reaping it kills
the interpreter that is doing it. Selection and removal happen in the
same critical section as :meth:_begin_request, so a kernel cannot be
chosen here and handed a cell in the same instant.
Source code in tit/server/kernels.py
| def reap_idle(self) -> list[str]:
"""Shut down kernels nobody has used lately. Returns their ids.
*Idle* means no request in flight and not executing — a cell that runs
for longer than the timeout is work, not neglect, and reaping it kills
the interpreter that is doing it. Selection and removal happen in the
same critical section as :meth:`_begin_request`, so a kernel cannot be
chosen here and handed a cell in the same instant.
"""
if self.idle_timeout <= 0:
return []
cutoff = self._clock() - self.idle_timeout
with self._lock:
stale = [
s
for s in self._sessions.values()
if s.in_flight == 0
and s.execution_state not in ("busy", "starting")
and s.last_used < cutoff
]
for session in stale:
self._sessions.pop(session.id, None)
for session in stale:
logger.info("kernel %s reaped after %.0fs idle", session.id, self.idle_timeout)
self._shutdown_session(session)
return [s.id for s in stale]
|
complete
complete(kernel_id: str, req_id: str, code: str, cursor_pos: int) -> None
Ask the kernel what could follow the cursor.
This is complete_request, the same call Jupyter's own front end
makes. It is answered by IPython's completer against the kernel's
live namespace, so after the notebook has run its imports,
tit.<Tab> lists what tit actually holds in that interpreter.
Source code in tit/server/kernels.py
| def complete(self, kernel_id: str, req_id: str, code: str, cursor_pos: int) -> None:
"""Ask the kernel what could follow the cursor.
This is `complete_request`, the same call Jupyter's own front end
makes. It is answered by IPython's completer against the kernel's
**live namespace**, so after the notebook has run its imports,
``tit.<Tab>`` lists what `tit` actually holds in that interpreter.
"""
session = self.get(kernel_id)
self._begin_request(session)
try:
msg_id = session.client.complete(code, cursor_pos)
except Exception as error:
self._end_request(session)
raise KernelError("op-failed", f"complete: {error}") from error
with session.lock:
session.queries[msg_id] = (req_id, "complete")
|
inspect
inspect(kernel_id: str, req_id: str, code: str, cursor_pos: int, detail: int = 0) -> None
Ask the kernel about the name under the cursor (Jupyter's ⇧⇥).
Source code in tit/server/kernels.py
| def inspect(self, kernel_id: str, req_id: str, code: str, cursor_pos: int, detail: int = 0) -> None:
"""Ask the kernel about the name under the cursor (Jupyter's ⇧⇥)."""
session = self.get(kernel_id)
self._begin_request(session)
try:
msg_id = session.client.inspect(code, cursor_pos, detail_level=detail)
except Exception as error:
self._end_request(session)
raise KernelError("op-failed", f"inspect: {error}") from error
with session.lock:
session.queries[msg_id] = (req_id, "inspect")
|
restart
restart(kernel_id: str) -> None
Restart the interpreter, rebuilding the client around it.
The pumps are stopped and joined first (:meth:_stop_pumps says why),
and the client is REPLACED rather than reused: restart_kernel gives
the new interpreter a new session key, and the old client's channels
then reject every message with Invalid Signature -- which is what
they did, until the socket teardown underneath aborted the process.
Source code in tit/server/kernels.py
| def restart(self, kernel_id: str) -> None:
"""Restart the interpreter, rebuilding the client around it.
The pumps are stopped and joined first (:meth:`_stop_pumps` says why),
and the client is REPLACED rather than reused: ``restart_kernel`` gives
the new interpreter a new session key, and the old client's channels
then reject every message with ``Invalid Signature`` -- which is what
they did, until the socket teardown underneath aborted the process.
"""
session = self.get(kernel_id)
session.last_used = self._clock()
with self._lock:
# Nothing survives a restart, so nothing is in flight afterwards.
session.in_flight = 0
with session.lock:
session.pending.clear()
session.finishing.clear()
session.queries.clear()
self._emit(session, {"type": "status", "state": "starting"})
self._stop_pumps(session)
try:
session.client.stop_channels()
except Exception as error: # pragma: no cover - already-closed channels
logger.debug("kernel %s stop_channels before restart: %s", session.id, error)
try:
session.manager.restart_kernel(now=False)
client = session.manager.client()
client.start_channels()
client.wait_for_ready(timeout=STARTUP_TIMEOUT_SECONDS)
except Exception as error:
# The pumps stay down: there is no live channel for them to poll,
# and starting them on a dead kernel is how the abort happened.
self._emit(session, {"type": "status", "state": "dead"})
raise KernelError("op-failed", f"restart: {error}") from error
session.client = client
session.execution_state = "idle"
self._start_pumps(session)
self._emit(session, {"type": "ready", "kernel": session.describe()})
|
subscribe
subscribe(kernel_id: str, listener: Listener) -> Callable[[], None]
Attach a listener; returns the function that detaches it.
Source code in tit/server/kernels.py
| def subscribe(self, kernel_id: str, listener: Listener) -> Callable[[], None]:
"""Attach a listener; returns the function that detaches it."""
session = self.get(kernel_id)
with session.lock:
session.listeners.append(listener)
def unsubscribe() -> None:
with session.lock:
if listener in session.listeners:
session.listeners.remove(listener)
return unsubscribe
|
output_from_msg
An iopub message as the nbformat output it will be stored as.
Source code in tit/server/kernels.py
| def output_from_msg(msg_type: str, content: dict[str, Any]) -> dict[str, Any] | None:
"""An iopub message as the nbformat output it will be stored as."""
if msg_type == "stream":
return {
"output_type": "stream",
"name": content.get("name", "stdout"),
"text": content.get("text", ""),
}
if msg_type in ("display_data", "update_display_data"):
return {
"output_type": "display_data",
"data": content.get("data", {}),
"metadata": content.get("metadata", {}),
}
if msg_type == "execute_result":
return {
"output_type": "execute_result",
"data": content.get("data", {}),
"metadata": content.get("metadata", {}),
"execution_count": content.get("execution_count"),
}
if msg_type == "error":
return {
"output_type": "error",
"ename": content.get("ename", ""),
"evalue": content.get("evalue", ""),
"traceback": content.get("traceback", []),
}
return None
|
get_kernel_registry
The process-wide registry, created on first use.
Source code in tit/server/kernels.py
| def get_kernel_registry() -> KernelRegistry:
"""The process-wide registry, created on first use."""
global _registry
if _registry is None:
_registry = KernelRegistry()
return _registry
|
reset_kernel_registry
reset_kernel_registry() -> None
Drop the registry, shutting down whatever it still holds. Tests.
Source code in tit/server/kernels.py
| def reset_kernel_registry() -> None:
"""Drop the registry, shutting down whatever it still holds. Tests."""
global _registry
if _registry is not None:
_registry.shutdown_all()
_registry = None
|