feat(coding-agent/eval): added local python-runner subprocess execution

- Replaced Python execution with a local `python -u runner.py` subprocess and NDJSON stdin/stdout framing.
- Removed shared-gateway architecture, including coordinator lifecycle APIs, `useSharedGateway` wiring, and `jupyter` CLI/actions.
- Simplified setup checks to a plain Python 3 availability probe and removed automatic dependency-install fallbacks.
- Updated kernel cancellation and display processing to use status frames, SIGINT/SIGTERM escalation, and normalized output coercion.
- Added `python-runner` integration and display tests while deleting legacy websocket and kernel lifecycle test suites.
This commit is contained in:
can1357
2026-05-12 09:09:24 +02:00
parent dcf3f1c8a3
commit 8d144e17ec
26 changed files with 1735 additions and 3004 deletions
+3 -4
View File
@@ -252,10 +252,9 @@ Related vars:
| Variable | Default / behavior | | Variable | Default / behavior |
| ------------------------- | ------------------------------------------------------------------------------------------------------------------- | | ------------------------- | ------------------------------------------------------------------------------------------------------------------- |
| `PI_PY` | Eval backend override: `0`/`bash`=JavaScript only, `1`/`py`=Python only, `mix`/`both`=both; invalid values ignored | | `PI_PY` | Eval backend override: `0`/`bash`=JavaScript only, `1`/`py`=Python only, `mix`/`both`=both; invalid values ignored |
| `PI_PYTHON_SKIP_CHECK` | If `1`, skips Python kernel availability checks/warm checks | | `PI_PYTHON_SKIP_CHECK` | If `1`, skips Python interpreter availability checks (subprocess runner still starts on demand) |
| `PI_PYTHON_GATEWAY_URL` | If set, uses external kernel gateway instead of local shared gateway | | `PI_PYTHON_INTEGRATION` | If `1`, opts gated integration tests in (e.g. `python-runner.integration.test.ts`) into running against real Python |
| `PI_PYTHON_GATEWAY_TOKEN` | Optional auth token for external gateway (`Authorization: token <value>`) | | `PI_PYTHON_IPC_TRACE` | If `1`, logs NDJSON frames exchanged with the Python runner subprocess |
| `PI_PYTHON_IPC_TRACE` | If `1`, enables low-level IPC trace path in kernel module |
| `VIRTUAL_ENV` | Highest-priority venv path for Python runtime resolution | | `VIRTUAL_ENV` | Highest-priority venv path for Python runtime resolution |
Extra conditional behavior: Extra conditional behavior:
+105 -158
View File
@@ -1,20 +1,22 @@
# Eval Tool Python Backend and IPython Runtime # Eval Tool Python Backend
This document describes the current Python execution stack in `packages/coding-agent`. This document describes the Python execution stack in `packages/coding-agent`.
It covers tool behavior, kernel/gateway lifecycle, environment handling, execution semantics, output rendering, and operational failure modes. It covers tool behavior, runner lifecycle, environment handling, execution semantics, output rendering, supported magics, and operational failure modes.
## Scope and Key Files ## Scope and Key Files
- Tool surface: `src/tools/eval.ts` - Tool surface: `src/tools/eval.ts`
- Session/per-call kernel orchestration: `src/eval/py/executor.ts` - Session/per-call kernel orchestration: `src/eval/py/executor.ts`
- Kernel protocol + gateway integration: `src/eval/py/kernel.ts` - Subprocess kernel client: `src/eval/py/kernel.ts`
- Shared local gateway coordinator: `src/eval/py/gateway-coordinator.ts` - Python wrapper / NDJSON server: `src/eval/py/runner.py`
- Prelude helpers loaded into every kernel: `src/eval/py/prelude.py`
- MIME bundle renderer (text + structured outputs): `src/eval/py/display.ts`
- Interactive-mode renderer for user-triggered Python runs: `src/modes/components/eval-execution.ts` - Interactive-mode renderer for user-triggered Python runs: `src/modes/components/eval-execution.ts`
- Runtime/env filtering and Python resolution: `src/eval/py/runtime.ts` - Runtime/env filtering and Python resolution: `src/eval/py/runtime.ts`
## What eval's Python backend is ## What eval's Python backend is
The `eval` tool executes one or more Python cells through a Jupyter Kernel Gateway-backed kernel when `language: "python"` is selected or inferred (not by spawning `python -c` directly per cell). The `eval` tool executes one or more Python cells inside a long-lived `python3` subprocess that speaks NDJSON over stdin/stdout. No Jupyter, no kernel gateway, no extra pip dependencies — a vanilla Python 3.8+ interpreter is enough. Rich `display()` output (PIL, pandas, plotly, matplotlib figures) keeps working because the wrapper reimplements the MIME-bundle dispatch that IPython previously provided.
Tool params: Tool params:
@@ -28,66 +30,73 @@ Tool params:
The tool is `concurrency = "exclusive"` for a session, so calls do not overlap. The tool is `concurrency = "exclusive"` for a session, so calls do not overlap.
## Gateway lifecycle
### Modes
There are two gateway paths:
1. **External gateway** (`PI_PYTHON_GATEWAY_URL` set)
- Uses the configured URL directly.
- Optional auth with `PI_PYTHON_GATEWAY_TOKEN`.
- No local gateway process is spawned or managed.
2. **Local shared gateway** (default path)
- Uses a single shared process coordinated under `~/.omp/agent/python-gateway`.
- Metadata file: `gateway.json`
- Lock file: `gateway.lock`
- Spawn command:
- `python -m kernel_gateway`
- bound to `127.0.0.1:<allocated-port>`
- startup health check: `GET /api/kernelspecs`
### Local shared gateway coordination
`acquireSharedGateway()`:
- Takes a file lock (`gateway.lock`) with heartbeat.
- Reuses `gateway.json` if PID is alive and health check passes.
- Cleans stale info/PIDs when needed.
- Starts a new gateway when no healthy one exists.
`releaseSharedGateway()` is currently a no-op (kernel shutdown does not tear down shared gateway).
`shutdownSharedGateway()` explicitly terminates the shared process and clears gateway metadata.
### Important constraint
`python.sharedGateway=false` is rejected at kernel start:
- Error: `Shared Python gateway required; local gateways are disabled`
- There is no per-process non-shared local gateway mode.
## Kernel lifecycle ## Kernel lifecycle
Kernels are created via `POST /api/kernels` on the selected gateway when a retained session needs a kernel or when `per-call` mode starts a request. Each kernel is a single Python subprocess: `python -u <runner.py>`. The runner is bundled with the host binary (Bun text import), written to `~/.omp/python-env`-adjacent tmp cache once per script-hash, and reused by every subsequent spawn.
Kernel startup sequence: Kernel startup sequence:
1. Availability check (`checkPythonKernelAvailability`) 1. Availability check (`checkPythonKernelAvailability`) — verifies that a Python interpreter resolves and runs.
2. Create kernel (`/api/kernels`) 2. Spawn `python -u runner.py` with filtered env and `cwd`.
3. Open websocket (`/api/kernels/:id/channels`) 3. Send an init request that runs `os.chdir(cwd)`, injects env entries, and adds `cwd` to `sys.path`.
4. Initialize kernel env (`cwd`, env vars, `sys.path`) 4. Execute `PYTHON_PRELUDE` (idempotent — only initializes once per process).
5. Execute `PYTHON_PRELUDE`
6. Load extension modules from:
- user: `~/.omp/agent/modules/*.py`
- project: `<cwd>/.omp/modules/*.py` (overrides same-name user module)
Kernel shutdown: Kernel shutdown:
- Deletes remote kernel via `DELETE /api/kernels/:id` - Send `{"type": "exit"}` over stdin.
- Closes websocket - Wait for process exit with `SHUTDOWN_GRACE_MS` budget.
- Calls shared gateway release hook (no-op today) - Escalate to `SIGTERM` and finally `SIGKILL` if the process does not exit in time.
## Wire protocol (NDJSON, host ↔ runner)
One JSON object per line, UTF-8, `\n` terminated.
Host → runner:
```jsonc
{"id": "<reqId>", "code": "<source>", "silent": false, "storeHistory": true}
{"type": "exit"}
```
Runner → host:
```jsonc
{"type": "started", "id": "<reqId>"}
{"type": "stdout", "id": "<reqId>", "data": "..."}
{"type": "stderr", "id": "<reqId>", "data": "..."}
{"type": "display", "id": "<reqId>", "bundle": {<mime>: <value>}}
{"type": "result", "id": "<reqId>", "bundle": {<mime>: <value>}}
{"type": "error", "id": "<reqId>", "ename": "...", "evalue": "...", "traceback": ["..."]}
{"type": "done", "id": "<reqId>", "status": "ok"|"error", "executionCount": N, "cancelled": false}
```
Status events the prelude emits (e.g. `_emit_status("find", count=…)`) ship inside display bundles under `application/x-omp-status` so the existing TUI status renderer keeps working.
## Magics
The runner's source transformer rewrites IPython-style magics to plain Python calls before parsing. Supported set:
| Magic | Effect |
| --- | --- |
| `%pip <args>` | `python -m pip <args>` with live streaming output. Newly installed packages are evicted from `sys.modules` so the next `import` picks up the fresh install. |
| `%cd <path>` | `os.chdir(path)` (with `~` expansion); emits status event. |
| `%pwd` | Returns `os.getcwd()`. |
| `%ls [path]` | Returns `sorted(os.listdir(path))`. |
| `%env [KEY[=VAL]]` | List, read, or set env vars (matches prelude `env()` semantics). |
| `%set_env KEY VALUE` | Set `os.environ[KEY]`. |
| `%time <expr>` / `%timeit <expr>` | Time the expression; emits status event with elapsed ms. |
| `%who` / `%whos` | List user-namespace names. |
| `%reset` | Clear user globals and re-inject prelude. |
| `%load <path>` | Read a file into a fresh cell and execute. |
| `%run <path>` | `runpy.run_path` and merge globals back. |
| `%%bash` / `%%sh` | Run the cell body via `bash`/`sh`. |
| `%%capture [name]` | Run body with stdout/stderr captured into `name`. |
| `%%timeit` | Time the cell body. |
| `%%writefile <path>` | Write body to file. |
| `!cmd` / `var = !cmd` | Run command via subprocess shell; returns an SList-style result with `.n` / `.s` helpers. |
| `var = %name args` | Assignment forms work for line magics and `!cmd`. |
Unknown magic names raise `NameError: UsageError: ...` inside the cell.
## Session persistence semantics ## Session persistence semantics
@@ -99,11 +108,10 @@ Kernel shutdown:
- Idle sessions are evicted after 5 minutes. - Idle sessions are evicted after 5 minutes.
- At most 4 sessions; oldest is evicted on overflow. - At most 4 sessions; oldest is evicted on overflow.
- Heartbeat checks detect dead kernels. - Heartbeat checks detect dead kernels.
- Auto-restart allowed once; repeated crash => hard failure. - Auto-restart allowed once; repeated crash ⇒ hard failure.
- `per-call` - `per-call`
- Creates a fresh kernel for each execute request. - Spawns a fresh subprocess for each request.
- Shuts kernel down after the request. - Shuts the subprocess down after the request.
- No cross-call state persistence. - No cross-call state persistence.
### Multi-cell behavior in a single tool call ### Multi-cell behavior in a single tool call
@@ -120,7 +128,7 @@ If an intermediate cell fails:
## Environment filtering and runtime resolution ## Environment filtering and runtime resolution
Environment is filtered before launching gateway/kernel runtime: Environment is filtered before launching the runner:
- Allowlist includes core vars like `PATH`, `HOME`, locale vars, `VIRTUAL_ENV`, `PYTHONPATH`, etc. - Allowlist includes core vars like `PATH`, `HOME`, locale vars, `VIRTUAL_ENV`, `PYTHONPATH`, etc.
- Allow-prefixes: `LC_`, `XDG_`, `PI_` - Allow-prefixes: `LC_`, `XDG_`, `PI_`
@@ -134,15 +142,7 @@ Runtime selection order:
When a venv is selected, its bin/Scripts path is prepended to `PATH`. When a venv is selected, its bin/Scripts path is prepended to `PATH`.
Kernel startup receives the optional session file path from the executor: The runner additionally receives `PYTHONUNBUFFERED=1` and `PYTHONIOENCODING=utf-8` so streamed output reaches the host promptly.
- `PI_SESSION_FILE` (session state file path)
`PythonKernel.#initializeKernelEnvironment(...)` then runs init script inside kernel to:
- `os.chdir(cwd)`
- injects provided env entries into `os.environ`
- ensures cwd is in `sys.path`
## Tool availability and mode selection ## Tool availability and mode selection
@@ -154,9 +154,9 @@ Kernel startup receives the optional session file path from the executor:
`PI_PY` accepted values: `PI_PY` accepted values:
- `0` / `bash` -> JavaScript backend only - `0` / `bash` → JavaScript backend only
- `1` / `py` -> Python backend only - `1` / `py` → Python backend only
- `mix` / `both` -> both backends - `mix` / `both` → both backends
If Python preflight fails and `eval.js` is enabled, `eval` remains available and dispatches to JavaScript unless `language: "python"` is explicitly requested. If Python preflight fails and `eval.js` is enabled, `eval` remains available and dispatches to JavaScript unless `language: "python"` is explicitly requested.
@@ -164,45 +164,33 @@ If Python preflight fails and `eval.js` is enabled, `eval` remains available and
### Tool-level timeout ### Tool-level timeout
`eval` timeout is in seconds, default 30, clamped to `1..600`. `eval` timeout is in seconds, default 30, clamped to `1..600`. The tool combines caller abort signal and timeout signal with `AbortSignal.any(...)`.
The tool combines:
- caller abort signal
- timeout abort signal
with `AbortSignal.any(...)`.
### Kernel execution cancellation ### Kernel execution cancellation
On abort/timeout: On abort/timeout:
- Execution is marked cancelled. - The host sends `kill("SIGINT")` to the runner subprocess.
- Kernel interrupt is attempted via REST (`POST /interrupt`) and control-channel `interrupt_request`. - The runner's exec-time signal handler raises `KeyboardInterrupt` inside the user code.
- Result includes `cancelled=true`. - Result includes `cancelled=true`; timeout path annotates output as `Command timed out after <n> seconds`.
- Timeout path annotates output as `Command timed out after <n> seconds`. - Between requests the runner installs `SIG_IGN` for SIGINT so a stray cancel does not tear down the kernel.
If a second cancel is required (runner stuck in C code), the host escalates to `SIGTERM` and the session restarts on the next call.
### stdin behavior ### stdin behavior
Interactive stdin is not supported. Interactive stdin is not supported. The runner does not forward `input()` prompts; user code that calls `input()` blocks until cancellation.
If kernel emits `input_request`:
- Tool records `stdinRequested=true`
- Emits explanatory text
- Sends empty `input_reply`
- Execution is treated as failure at executor layer
## Output capture and rendering ## Output capture and rendering
### Captured output classes ### Captured output classes
From kernel messages: From runner frames:
- `stream` -> plain text chunks - `stdout` / `stderr` → plain text chunks
- `display_data`/`execute_result` -> rich display handling - `display` / `result` → rich display handling (MIME bundle)
- `error` -> traceback text - `error` → traceback text
- custom MIME `application/x-omp-status` -> structured status events - `application/x-omp-status` MIME inside `display` → structured status events
Display MIME precedence: Display MIME precedence:
@@ -212,15 +200,17 @@ Display MIME precedence:
Additionally captured as structured outputs: Additionally captured as structured outputs:
- `application/json` -> JSON tree data - `application/json` → JSON tree data
- `image/png` -> image payloads - `image/png` / `image/jpeg` → image payloads
- `application/x-omp-status` -> status events - `application/x-omp-status` → status events
### Matplotlib
The runner sets `MPLBACKEND=Agg` as an environ default so figures render off-screen. After every cell, `pyplot.get_fignums()` is iterated; each figure is saved to PNG, emitted as an `image/png` display, and closed.
### Storage and truncation ### Storage and truncation
Output is streamed through `OutputSink` and may be persisted to artifact storage. Output is streamed through `OutputSink` and may be persisted to artifact storage. Tool results can include truncation metadata and `artifact://<id>` for full output recovery.
Tool results can include truncation metadata and `artifact://<id>` for full output recovery.
### Renderer behavior ### Renderer behavior
@@ -234,61 +224,18 @@ Tool results can include truncation metadata and `artifact://<id>` for full outp
- clamps very long individual lines to 4000 chars for display safety - clamps very long individual lines to 4000 chars for display safety
- shows cancellation/error/truncation notices - shows cancellation/error/truncation notices
## External gateway support ## Operational troubleshooting
Set: - **Python backend not available** — Check `eval.py`, `PI_PY`, and that `python`/`python3` is on PATH. If preflight fails and `eval.js` is enabled, omit `language` or pass `language: "js"` to use JavaScript.
- **No Python on PATH** — Install a system Python 3.8+ or place a venv at `~/.omp/python-env`. `omp setup python --check` reports the resolved interpreter.
```bash - **Execution hangs then times out** — Increase tool `timeout` (max 600s) if workload is legitimate. For stuck native code, cancellation triggers `SIGINT` first then escalates; the session restarts on the next request.
export PI_PYTHON_GATEWAY_URL="http://127.0.0.1:8888" - **stdin/input prompts in Python code** — `input()` is not supported; pass data programmatically.
# Optional: - **Working directory errors** — Tool validates `cwd` exists and is a directory before execution.
export PI_PYTHON_GATEWAY_TOKEN="..."
```
Behavior differences from local shared gateway:
- No local gateway lock/info files
- No local process spawn/termination
- Health checks and kernel CRUD run against external endpoint
- Auth failures are surfaced with explicit token guidance
## Operational troubleshooting (current failure modes)
- **Python backend not available**
- Check `eval.py`, Python kernel dependencies, and `PI_PY`.
- If preflight fails and `eval.js` is enabled, omit `language` or pass `language: "js"` to use JavaScript.
- **Kernel availability errors**
- Local mode requires both `kernel_gateway` and `ipykernel` importable in resolved Python runtime.
- Install with:
```bash
python -m pip install jupyter_kernel_gateway ipykernel
```
- **`python.sharedGateway=false` causes startup failure**
- This is expected with current implementation.
- **External gateway auth/reachability failures**
- 401/403 -> set `PI_PYTHON_GATEWAY_TOKEN`.
- timeout/unreachable -> verify URL/network and gateway health.
- **Execution hangs then times out**
- Increase tool `timeout` (max 600s) if workload is legitimate.
- For stuck code, cancellation triggers kernel interrupt but user code may still need refactor.
- **stdin/input prompts in Python code**
- `input()` is not supported interactively in this runtime path; pass data programmatically.
- **Resource exhaustion (`EMFILE` / too many open files)**
- Session manager triggers shared-gateway recovery (session teardown + shared gateway restart).
- **Working directory errors**
- Tool validates `cwd` exists and is a directory before execution.
## Relevant environment variables ## Relevant environment variables
- `PI_PY` — tool exposure override (`bash-only`/`ipy-only`/`both` mapping above) - `PI_PY` — tool exposure override
- `PI_PYTHON_GATEWAY_URL` — use external gateway
- `PI_PYTHON_GATEWAY_TOKEN` — optional external gateway auth token
- `PI_PYTHON_SKIP_CHECK=1` — bypass Python preflight/warm checks - `PI_PYTHON_SKIP_CHECK=1` — bypass Python preflight/warm checks
- `PI_PYTHON_IPC_TRACE=1` — log kernel IPC send/receive traces - `PI_PYTHON_INTEGRATION=1` — enable gated integration tests that spawn a real Python
- `PI_PYTHON_IPC_TRACE=1` — log NDJSON frames exchanged with the runner subprocess
- `PI_DEBUG_STARTUP=1` — emit startup-stage debug markers - `PI_DEBUG_STARTUP=1` — emit startup-stage debug markers
+2 -3
View File
@@ -177,12 +177,11 @@ A single tool call can mix Python and JS cells. Persistence is per language runt
- JS/Python prelude helpers can read, write, append, diff, and traverse files under the session cwd or absolute paths. - JS/Python prelude helpers can read, write, append, diff, and traverse files under the session cwd or absolute paths.
- Output may spill to an artifact file via `OutputSink`. - Output may spill to an artifact file via `OutputSink`.
- Network - Network
- Python backend talks to a Jupyter kernel gateway over HTTP and WebSocket. - Python backend speaks NDJSON to a local `python3` subprocess over stdin/stdout (no network).
- External gateway mode uses `PI_PYTHON_GATEWAY_URL` and optional `PI_PYTHON_GATEWAY_TOKEN`.
- JS runtime exposes `fetch` and `tool.<name>()`; those tools may perform additional network I/O. - JS runtime exposes `fetch` and `tool.<name>()`; those tools may perform additional network I/O.
- Subprocesses / native bindings - Subprocesses / native bindings
- Python availability check runs `<python> -c ...`. - Python availability check runs `<python> -c ...`.
- Python backend may start or connect to a kernel gateway; details are in `docs/python-repl.md`. - Python backend spawns one `python -u runner.py` subprocess per kernel; cancellation sends `SIGINT`. Details in `docs/python-repl.md`.
- Session state - Session state
- `session.assertEvalExecutionAllowed?.()` can block execution. - `session.assertEvalExecutionAllowed?.()` can block execution.
- `session.trackEvalExecution?.(...)` can register cancellable eval work. - `session.trackEvalExecution?.(...)` can register cancellable eval work.
+8 -1
View File
@@ -1,9 +1,11 @@
# Changelog # Changelog
## [Unreleased] ## [Unreleased]
### Breaking Changes ### Breaking Changes
- Changed the `timeoutMs` execution option to no longer be enforced during worker-based JS runs, so callers must rely on external cancellation signals for time limits - Changed the `timeoutMs` execution option to no longer be enforced during worker-based JS runs, so callers must rely on external cancellation signals for time limits
- Replaced the Jupyter kernel gateway + WebSocket protocol behind the Python `eval` backend with a subprocess-backed runner that speaks NDJSON over stdin/stdout; removed the `jupyter_kernel_gateway`/`ipykernel` pip dependencies, the `python.sharedGateway` setting, the `omp jupyter` CLI command, and the `PI_PYTHON_GATEWAY_URL` / `PI_PYTHON_GATEWAY_TOKEN` environment variables
### Added ### Added
@@ -18,11 +20,16 @@
### Changed ### Changed
- Changed JavaScript execution in `executeJs` to expose the worker’s real `process` object instead of a restricted, frozen subset - Changed `setup python` to only verify a reachable Python 3 interpreter instead of installing Jupyter dependencies
- Changed `info` output to remove the obsolete Python Gateway status block now that shared gateway management is no longer available
- Changed JavaScript execution in `executeJs` to expose the worker\u2019s real `process` object instead of a restricted, frozen subset
- Changed JavaScript evaluation to run per session in a worker-backed runner with explicit initialization and teardown handling - Changed JavaScript evaluation to run per session in a worker-backed runner with explicit initialization and teardown handling
- Changed the Python backend to launch one `python -u runner.py` subprocess per kernel; cancellation now sends `SIGINT` which raises a real `KeyboardInterrupt` in user code, and the same subprocess is reused across cells in session mode
- Changed Python magic handling so `%pip`, `%cd`, `%env`, `%pwd`, `%ls`, `%time`, `%timeit`, `%who`, `%reset`, `%load`, `%run`, `%%bash`, `%%capture`, `%%timeit`, `%%writefile`, and `!shell` work without depending on IPython
### Fixed ### Fixed
- Fixed Python output rendering so `text/markdown` takes precedence over `text/plain` and status bundles are emitted as status updates rather than plain text
- Fixed query tokenization in `HistoryStorage.search` so punctuation-delimited terms like `git-commit` are aligned with indexing and matched correctly - Fixed query tokenization in `HistoryStorage.search` so punctuation-delimited terms like `git-commit` are aligned with indexing and matched correctly
- Fixed history search result merging to de-duplicate matches and return full-text matches before substring-only matches while still respecting the requested limit - Fixed history search result merging to de-duplicate matches and return full-text matches before substring-only matches while still respecting the requested limit
- Fixed JS run cancellation so aborting a run now also cancels in-flight tool calls and terminates the active worker session - Fixed JS run cancellation so aborting a run now also cancels in-flight tool calls and terminates the active worker session
-1
View File
@@ -55,7 +55,6 @@ const commands: CommandEntry[] = [
{ name: "config", load: () => import("./commands/config").then(m => m.default) }, { name: "config", load: () => import("./commands/config").then(m => m.default) },
{ name: "grep", load: () => import("./commands/grep").then(m => m.default) }, { name: "grep", load: () => import("./commands/grep").then(m => m.default) },
{ name: "grievances", load: () => import("./commands/grievances").then(m => m.default) }, { name: "grievances", load: () => import("./commands/grievances").then(m => m.default) },
{ name: "jupyter", load: () => import("./commands/jupyter").then(m => m.default) },
{ name: "plugin", load: () => import("./commands/plugin").then(m => m.default) }, { name: "plugin", load: () => import("./commands/plugin").then(m => m.default) },
{ name: "setup", load: () => import("./commands/setup").then(m => m.default) }, { name: "setup", load: () => import("./commands/setup").then(m => m.default) },
{ name: "shell", load: () => import("./commands/shell").then(m => m.default) }, { name: "shell", load: () => import("./commands/shell").then(m => m.default) },
@@ -1,106 +0,0 @@
/**
* Jupyter CLI command handlers.
*
* Handles `omp jupyter` subcommand for managing the shared Python gateway.
*/
import { APP_NAME } from "@oh-my-pi/pi-utils";
import chalk from "chalk";
import { getGatewayStatus, shutdownSharedGateway } from "../eval/py/gateway-coordinator";
export type JupyterAction = "kill" | "status";
export interface JupyterCommandArgs {
action: JupyterAction;
}
export function parseJupyterArgs(args: string[]): JupyterCommandArgs | undefined {
if (args.length === 0 || args[0] !== "jupyter") {
return undefined;
}
const action = args[1] as JupyterAction | undefined;
if (!action || !["kill", "status"].includes(action)) {
return { action: "status" };
}
return { action };
}
export async function runJupyterCommand(cmd: JupyterCommandArgs): Promise<void> {
switch (cmd.action) {
case "kill":
await runKill();
break;
case "status":
await runStatus();
break;
}
}
async function runKill(): Promise<void> {
const status = await getGatewayStatus();
if (!status.active) {
console.log(chalk.dim("No Jupyter gateway is running"));
return;
}
console.log(`Killing Jupyter gateway (PID ${status.pid})...`);
await shutdownSharedGateway();
console.log(chalk.green("Jupyter gateway stopped"));
}
async function runStatus(): Promise<void> {
const status = await getGatewayStatus();
if (!status.active) {
console.log(chalk.dim("No Jupyter gateway is running"));
return;
}
console.log(chalk.bold("Jupyter Gateway Status\n"));
console.log(` ${chalk.green("●")} Running`);
console.log(` PID: ${status.pid}`);
console.log(` URL: ${status.url}`);
if (status.uptime !== null) {
console.log(` Uptime: ${formatUptime(status.uptime)}`);
}
if (status.pythonPath) {
console.log(` Python: ${status.pythonPath}`);
}
if (status.venvPath) {
console.log(` Venv: ${status.venvPath}`);
}
}
function formatUptime(ms: number): string {
const seconds = Math.floor(ms / 1000);
const minutes = Math.floor(seconds / 60);
const hours = Math.floor(minutes / 60);
if (hours > 0) {
return `${hours}h ${minutes % 60}m`;
}
if (minutes > 0) {
return `${minutes}m ${seconds % 60}s`;
}
return `${seconds}s`;
}
export function printJupyterHelp(): void {
console.log(`${chalk.bold(`${APP_NAME} jupyter`)} - Manage the shared Jupyter gateway
${chalk.bold("Usage:")}
${APP_NAME} jupyter <command>
${chalk.bold("Commands:")}
status Show gateway status (default)
kill Stop the running gateway
${chalk.bold("Examples:")}
${APP_NAME} jupyter # Show status
${APP_NAME} jupyter status # Show status
${APP_NAME} jupyter kill # Stop the gateway
`);
}
+14 -161
View File
@@ -21,7 +21,6 @@ export interface SetupCommandArgs {
const VALID_COMPONENTS: SetupComponent[] = ["python", "stt"]; const VALID_COMPONENTS: SetupComponent[] = ["python", "stt"];
const PYTHON_PACKAGES = ["jupyter_kernel_gateway", "ipykernel"];
const MANAGED_PYTHON_ENV = getPythonEnvDir(); const MANAGED_PYTHON_ENV = getPythonEnvDir();
/** /**
@@ -65,10 +64,6 @@ export function parseSetupArgs(args: string[]): SetupCommandArgs | undefined {
interface PythonCheckResult { interface PythonCheckResult {
available: boolean; available: boolean;
pythonPath?: string; pythonPath?: string;
uvPath?: string;
pipPath?: string;
missingPackages: string[];
installedPackages: string[];
usingManagedEnv?: boolean; usingManagedEnv?: boolean;
managedEnvPath?: string; managedEnvPath?: string;
} }
@@ -85,8 +80,6 @@ function managedPythonPath(): string {
async function checkPythonSetup(): Promise<PythonCheckResult> { async function checkPythonSetup(): Promise<PythonCheckResult> {
const result: PythonCheckResult = { const result: PythonCheckResult = {
available: false, available: false,
missingPackages: [],
installedPackages: [],
managedEnvPath: MANAGED_PYTHON_ENV, managedEnvPath: MANAGED_PYTHON_ENV,
}; };
@@ -94,109 +87,24 @@ async function checkPythonSetup(): Promise<PythonCheckResult> {
const managedPath = managedPythonPath(); const managedPath = managedPythonPath();
const hasManagedEnv = await Bun.file(managedPath).exists(); const hasManagedEnv = await Bun.file(managedPath).exists();
result.uvPath = $which("uv") ?? undefined; const pythonPath = systemPythonPath ?? (hasManagedEnv ? managedPath : undefined);
result.pipPath = $which("pip3") ?? $which("pip") ?? undefined; if (!pythonPath) {
const candidates = [systemPythonPath, hasManagedEnv ? managedPath : undefined].filter(
(candidate): candidate is string => !!candidate,
);
if (candidates.length === 0) {
return result; return result;
} }
const probe = await $`${pythonPath} -c "import sys;sys.exit(0)"`.quiet().nothrow();
result.pythonPath = systemPythonPath ?? managedPath; result.pythonPath = pythonPath;
let bestMatch = { result.available = probe.exitCode === 0;
pythonPath: candidates[0], result.usingManagedEnv = pythonPath === managedPath;
missingPackages: [...PYTHON_PACKAGES],
installedPackages: [] as string[],
usingManagedEnv: candidates[0] === managedPath,
};
for (const pythonPath of candidates) {
const installedPackages: string[] = [];
const missingPackages: string[] = [];
for (const pkg of PYTHON_PACKAGES) {
const moduleName = pkg === "jupyter_kernel_gateway" ? "kernel_gateway" : pkg;
const script = `import importlib.util; raise SystemExit(0 if importlib.util.find_spec('${moduleName}') else 1)`;
const check = await $`${pythonPath} -c ${script}`.quiet().nothrow();
if (check.exitCode === 0) {
installedPackages.push(pkg);
} else {
missingPackages.push(pkg);
}
}
if (missingPackages.length < bestMatch.missingPackages.length) {
bestMatch = {
pythonPath,
missingPackages,
installedPackages,
usingManagedEnv: pythonPath === managedPath,
};
}
if (missingPackages.length === 0) {
result.available = true;
result.pythonPath = pythonPath;
result.missingPackages = missingPackages;
result.installedPackages = installedPackages;
result.usingManagedEnv = pythonPath === managedPath;
return result;
}
}
result.pythonPath = bestMatch.pythonPath;
result.missingPackages = bestMatch.missingPackages;
result.installedPackages = bestMatch.installedPackages;
result.usingManagedEnv = bestMatch.usingManagedEnv;
return result; return result;
} }
/** /**
* Install Python packages using uv (preferred) or pip. * Install Python packages using uv (preferred) or pip.
*/ */
async function installPythonPackages( // Python installation helper removed: the subprocess runner has no Python
packages: string[], // package dependencies beyond a working interpreter. `omp setup python --check`
pythonPath: string, // remains as a probe; users install optional libs (pandas, matplotlib, ...)
uvPath?: string, // directly via pip or the in-process `%pip` magic.
pipPath?: string,
): Promise<{ success: boolean; usedManagedEnv: boolean }> {
if (uvPath) {
console.log(chalk.dim(`Installing via uv: ${packages.join(" ")}`));
const result = await $`${uvPath} pip install ${packages}`.nothrow();
if (result.exitCode === 0) {
return { success: true, usedManagedEnv: false };
}
}
if (pipPath) {
console.log(chalk.dim(`Installing via pip: ${packages.join(" ")}`));
const result = await $`${pipPath} install ${packages}`.nothrow();
if (result.exitCode === 0) {
return { success: true, usedManagedEnv: false };
}
}
console.log(chalk.dim(`Falling back to managed virtual environment: ${MANAGED_PYTHON_ENV}`));
if (uvPath) {
const createEnv = await $`${uvPath} venv ${MANAGED_PYTHON_ENV}`.quiet().nothrow();
if (createEnv.exitCode !== 0) {
return { success: false, usedManagedEnv: true };
}
const installInManagedEnv = await $`${uvPath} pip install --python ${MANAGED_PYTHON_ENV} ${packages}`.nothrow();
return { success: installInManagedEnv.exitCode === 0, usedManagedEnv: true };
}
const createEnv = await $`${pythonPath} -m venv ${MANAGED_PYTHON_ENV}`.quiet().nothrow();
if (createEnv.exitCode !== 0) {
return { success: false, usedManagedEnv: true };
}
const managedPython = managedPythonPath();
const installInManagedEnv = await $`${managedPython} -m pip install ${packages}`.nothrow();
return { success: installInManagedEnv.exitCode === 0, usedManagedEnv: true };
}
/** /**
* Run the setup command. * Run the setup command.
@@ -232,67 +140,13 @@ async function handlePythonSetup(flags: { json?: boolean; check?: boolean }): Pr
console.log(chalk.dim(`Using managed environment: ${check.managedEnvPath}`)); console.log(chalk.dim(`Using managed environment: ${check.managedEnvPath}`));
} }
if (check.uvPath) { if (check.available) {
console.log(chalk.dim(`uv: ${check.uvPath}`));
} else if (check.pipPath) {
console.log(chalk.dim(`pip: ${check.pipPath}`));
}
if (check.installedPackages.length > 0) {
console.log(chalk.green(`${theme.status.success} Installed: ${check.installedPackages.join(", ")}`));
}
if (check.missingPackages.length === 0) {
console.log(chalk.green(`\n${theme.status.success} Python execution is ready`)); console.log(chalk.green(`\n${theme.status.success} Python execution is ready`));
return; return;
} }
console.log(chalk.yellow(`${theme.status.warning} Missing: ${check.missingPackages.join(", ")}`)); console.error(chalk.red(`\n${theme.status.error} Python interpreter reported failure`));
process.exit(1);
if (flags.check) {
process.exit(1);
}
if (!check.uvPath && !check.pipPath) {
console.error(chalk.red(`\n${theme.status.error} No package manager found`));
console.error(chalk.dim("Install uv (recommended) or pip:"));
console.error(chalk.dim(" curl -LsSf https://astral.sh/uv/install.sh | sh"));
process.exit(1);
}
console.log("");
const install = await installPythonPackages(check.missingPackages, check.pythonPath, check.uvPath, check.pipPath);
if (!install.success) {
console.error(chalk.red(`\n${theme.status.error} Installation failed`));
console.error(chalk.dim("Try installing manually:"));
if (install.usedManagedEnv) {
if (check.uvPath) {
console.error(chalk.dim(` uv venv ${MANAGED_PYTHON_ENV}`));
console.error(
chalk.dim(` uv pip install --python ${MANAGED_PYTHON_ENV} ${check.missingPackages.join(" ")}`),
);
} else {
console.error(chalk.dim(` ${check.pythonPath} -m venv ${MANAGED_PYTHON_ENV}`));
console.error(chalk.dim(` ${managedPythonPath()} -m pip install ${check.missingPackages.join(" ")}`));
}
} else {
console.error(chalk.dim(` ${check.uvPath ? "uv pip" : "pip"} install ${check.missingPackages.join(" ")}`));
}
process.exit(1);
}
const recheck = await checkPythonSetup();
if (recheck.available) {
console.log(chalk.green(`\n${theme.status.success} Python execution is ready`));
if (recheck.usingManagedEnv) {
console.log(chalk.dim(`Managed Python environment: ${recheck.managedEnvPath}`));
}
} else {
console.error(chalk.red(`\n${theme.status.error} Setup incomplete`));
console.error(chalk.dim(`Still missing: ${recheck.missingPackages.join(", ")}`));
process.exit(1);
}
} }
async function handleSttSetup(flags: { json?: boolean; check?: boolean }): Promise<void> { async function handleSttSetup(flags: { json?: boolean; check?: boolean }): Promise<void> {
@@ -359,9 +213,8 @@ ${chalk.bold("Usage:")}
${APP_NAME} setup <component> [options] ${APP_NAME} setup <component> [options]
${chalk.bold("Components:")} ${chalk.bold("Components:")}
python Install Jupyter kernel dependencies for Python code execution python Verify a Python 3 interpreter is reachable for code execution
stt Install speech-to-text dependencies (openai-whisper, recording tools) stt Install speech-to-text dependencies (openai-whisper, recording tools)
Packages: ${PYTHON_PACKAGES.join(", ")}
${chalk.bold("Options:")} ${chalk.bold("Options:")}
-c, --check Check if dependencies are installed without installing -c, --check Check if dependencies are installed without installing
@@ -1,32 +0,0 @@
/**
* Manage the shared Jupyter gateway.
*/
import { Args, Command } from "@oh-my-pi/pi-utils/cli";
import { type JupyterAction, type JupyterCommandArgs, runJupyterCommand } from "../cli/jupyter-cli";
import { initTheme } from "../modes/theme/theme";
const ACTIONS: JupyterAction[] = ["kill", "status"];
export default class Jupyter extends Command {
static description = "Manage the shared Jupyter gateway";
static args = {
action: Args.string({
description: "Jupyter action",
required: false,
options: ACTIONS,
}),
};
async run(): Promise<void> {
const { args } = await this.parse(Jupyter);
const action = (args.action ?? "status") as JupyterAction;
const cmd: JupyterCommandArgs = {
action,
};
await initTheme();
await runJupyterCommand(cmd);
}
}
@@ -1662,16 +1662,6 @@ export const SETTINGS_SCHEMA = {
}, },
}, },
"python.sharedGateway": {
type: "boolean",
default: true,
ui: {
tab: "editing",
label: "Shared Python Gateway",
description: "Share IPython kernel gateway across pi instances",
},
},
// ──────────────────────────────────────────────────────────────────────── // ────────────────────────────────────────────────────────────────────────
// Tools // Tools
// ──────────────────────────────────────────────────────────────────────── // ────────────────────────────────────────────────────────────────────────
@@ -1,28 +0,0 @@
export function getAbortReason(signal: AbortSignal | undefined, fallbackReason: string): Error {
if (signal?.reason instanceof Error) return signal.reason;
if (typeof signal?.reason === "string" && signal.reason.length > 0) {
return new Error(signal.reason);
}
return new Error(fallbackReason);
}
export function createCancellationError(name: "AbortError" | "TimeoutError", message: string): Error {
const error = new Error(message);
error.name = name;
return error;
}
export function getExecutionCancellationError(
result: { timedOut?: boolean },
signal: AbortSignal | undefined,
fallbackReason: string,
): Error {
if (signal?.aborted) {
return getAbortReason(signal, fallbackReason);
}
if (result.timedOut) {
return createCancellationError("TimeoutError", fallbackReason);
}
return createCancellationError("AbortError", fallbackReason);
}
@@ -0,0 +1,71 @@
/**
* Display bundle rendering shared between the Python runner output and the
* legacy Jupyter MIME conventions. Pure function, no kernel coupling.
*/
import { htmlToBasicMarkdown } from "../../web/scrapers/types";
/** Status event emitted by prelude helpers for TUI rendering. */
export interface PythonStatusEvent {
/** Operation name (e.g., "find", "read", "write") */
op: string;
/** Additional data fields (count, path, pattern, etc.) */
[key: string]: unknown;
}
export type KernelDisplayOutput =
| { type: "json"; data: unknown }
| { type: "image"; data: string; mimeType: string }
| { type: "markdown" }
| { type: "status"; event: PythonStatusEvent };
function normalizeDisplayText(text: string): string {
return text.endsWith("\n") ? text : `${text}\n`;
}
/** Render a MIME bundle into text + structured outputs. */
export async function renderKernelDisplay(content: Record<string, unknown>): Promise<{
text: string;
outputs: KernelDisplayOutput[];
}> {
// Accept both raw bundles ({"text/plain": ...}) and Jupyter-style
// content envelopes ({ data: {...} }) so callers don't need to unwrap.
const data =
(content.data as Record<string, unknown> | undefined) ?? (content as Record<string, unknown> | undefined);
if (!data) return { text: "", outputs: [] };
const outputs: KernelDisplayOutput[] = [];
// Status events bypass the text path entirely — they exist only for TUI hooks.
if (data["application/x-omp-status"] !== undefined) {
const statusData = data["application/x-omp-status"];
if (statusData && typeof statusData === "object" && "op" in statusData) {
outputs.push({ type: "status", event: statusData as PythonStatusEvent });
}
return { text: "", outputs };
}
if (typeof data["image/png"] === "string") {
outputs.push({ type: "image", data: data["image/png"] as string, mimeType: "image/png" });
}
if (typeof data["image/jpeg"] === "string") {
outputs.push({ type: "image", data: data["image/jpeg"] as string, mimeType: "image/jpeg" });
}
if (data["application/json"] !== undefined) {
outputs.push({ type: "json", data: data["application/json"] });
}
// text/markdown takes precedence over text/plain (Markdown objects expose both
// where text/plain is just the repr).
if (typeof data["text/markdown"] === "string") {
outputs.push({ type: "markdown" });
return { text: normalizeDisplayText(String(data["text/markdown"])), outputs };
}
if (typeof data["text/plain"] === "string") {
return { text: normalizeDisplayText(String(data["text/plain"])), outputs };
}
if (data["text/html"] !== undefined) {
const markdown = (await htmlToBasicMarkdown(String(data["text/html"]))) || "";
return { text: markdown ? normalizeDisplayText(markdown) : "", outputs };
}
return { text: "", outputs };
}
+12 -88
View File
@@ -2,10 +2,9 @@ import { getProjectDir, logger } from "@oh-my-pi/pi-utils";
import { OutputSink } from "../../session/streaming-output"; import { OutputSink } from "../../session/streaming-output";
import type { ToolSession } from "../../tools"; import type { ToolSession } from "../../tools";
import type { JsStatusEvent } from "../js/shared/types"; import type { JsStatusEvent } from "../js/shared/types";
import { shutdownSharedGateway } from "./gateway-coordinator"; import type { KernelDisplayOutput } from "./display";
import { import {
checkPythonKernelAvailability, checkPythonKernelAvailability,
type KernelDisplayOutput,
type KernelExecuteOptions, type KernelExecuteOptions,
type KernelExecuteResult, type KernelExecuteResult,
PythonKernel, PythonKernel,
@@ -38,8 +37,6 @@ export interface PythonExecutorOptions {
kernelMode?: PythonKernelMode; kernelMode?: PythonKernelMode;
/** Restart the kernel before executing */ /** Restart the kernel before executing */
reset?: boolean; reset?: boolean;
/** Use shared gateway across pi instances (default: true) */
useSharedGateway?: boolean;
/** Session file path for accessing task outputs */ /** Session file path for accessing task outputs */
sessionFile?: string; sessionFile?: string;
/** /**
@@ -102,7 +99,6 @@ interface KernelSession {
restartCount: number; restartCount: number;
dead: boolean; dead: boolean;
needsRestart: boolean; needsRestart: boolean;
kernelInvalidatedByRecovery: boolean;
disposing: boolean; disposing: boolean;
disposeCapacityPromise?: Promise<void>; disposeCapacityPromise?: Promise<void>;
resolveDisposeCapacity?: () => void; resolveDisposeCapacity?: () => void;
@@ -122,7 +118,6 @@ const disposingKernelSessions = new Set<KernelSession>();
let cleanupTimer: NodeJS.Timeout | null = null; let cleanupTimer: NodeJS.Timeout | null = null;
interface KernelSessionExecutionOptions { interface KernelSessionExecutionOptions {
useSharedGateway?: boolean;
sessionFile?: string; sessionFile?: string;
artifactsDir?: string; artifactsDir?: string;
signal?: AbortSignal; signal?: AbortSignal;
@@ -295,7 +290,6 @@ function buildKernelStartOptions(
return { return {
cwd, cwd,
env, env,
useSharedGateway: options.useSharedGateway,
signal: options.signal, signal: options.signal,
deadlineMs: options.deadlineMs, deadlineMs: options.deadlineMs,
}; };
@@ -379,7 +373,6 @@ function finishDisposingKernelSession(session: KernelSession): void {
session.disposeResultPromise = undefined; session.disposeResultPromise = undefined;
session.disposeResultTimeoutMs = undefined; session.disposeResultTimeoutMs = undefined;
session.nextDisposalRetryAt = undefined; session.nextDisposalRetryAt = undefined;
session.kernelInvalidatedByRecovery = false;
syncCleanupTimer(); syncCleanupTimer();
} }
@@ -503,58 +496,6 @@ async function ensureKernelAvailable(
} }
} }
function isResourceExhaustionError(error: unknown): boolean {
const message = error instanceof Error ? error.message : String(error);
return (
message.includes("Too many open files") ||
message.includes("EMFILE") ||
message.includes("ENFILE") ||
message.includes("resource temporarily unavailable")
);
}
function clearSharedGatewayDisposingKernelSessionTracking(): void {
for (const session of Array.from(disposingKernelSessions.values())) {
if (!session.kernel.isSharedGateway) continue;
if (session.heartbeatTimer) {
clearInterval(session.heartbeatTimer);
session.heartbeatTimer = undefined;
}
disposingKernelSessions.delete(session);
session.resolveDisposeCapacity?.();
session.resolveDisposeCapacity = undefined;
session.disposeCapacityPromise = undefined;
session.resolveDisposeAttempt?.();
session.resolveDisposeAttempt = undefined;
session.disposeAttemptPromise = undefined;
session.disposeResultPromise = undefined;
session.disposeResultTimeoutMs = undefined;
session.nextDisposalRetryAt = undefined;
session.kernelInvalidatedByRecovery = false;
}
}
function markLiveKernelSessionsForRecovery(): void {
for (const session of kernelSessions.values()) {
if (session.heartbeatTimer) {
clearInterval(session.heartbeatTimer);
session.heartbeatTimer = undefined;
}
session.needsRestart = true;
session.kernelInvalidatedByRecovery = session.kernel.isSharedGateway;
session.restartCount = 0;
}
}
async function recoverFromResourceExhaustion(): Promise<void> {
logger.warn("Resource exhaustion detected, recovering by restarting shared gateway");
stopCleanupTimer();
markLiveKernelSessionsForRecovery();
clearSharedGatewayDisposingKernelSessionTracking();
await shutdownSharedGateway();
syncCleanupTimer();
}
function ensureKernelHeartbeat(session: KernelSession): void { function ensureKernelHeartbeat(session: KernelSession): void {
if (session.heartbeatTimer) return; if (session.heartbeatTimer) return;
session.heartbeatTimer = setInterval(() => { session.heartbeatTimer = setInterval(() => {
@@ -570,22 +511,12 @@ async function createKernelSession(
sessionId: string, sessionId: string,
cwd: string, cwd: string,
options: KernelSessionExecutionOptions = {}, options: KernelSessionExecutionOptions = {},
isRetry?: boolean,
): Promise<KernelSession> { ): Promise<KernelSession> {
requireRemainingTimeoutMs(options.deadlineMs); requireRemainingTimeoutMs(options.deadlineMs);
const env = buildKernelEnv(options); const env = buildKernelEnv(options);
const startOptions = buildKernelStartOptions(cwd, env, options); const startOptions = buildKernelStartOptions(cwd, env, options);
let kernel: PythonKernel; const kernel = await logger.time("createKernelSession:PythonKernel.start", PythonKernel.start, startOptions);
try {
kernel = await logger.time("createKernelSession:PythonKernel.start", PythonKernel.start, startOptions);
} catch (err) {
if (!isRetry && isResourceExhaustionError(err)) {
await recoverFromResourceExhaustion();
return createKernelSession(sessionId, cwd, options, true);
}
throw err;
}
const hasFallbackOwner = options.kernelOwnerId === undefined; const hasFallbackOwner = options.kernelOwnerId === undefined;
const initialOwnerId = options.kernelOwnerId ?? sessionId; const initialOwnerId = options.kernelOwnerId ?? sessionId;
@@ -596,7 +527,6 @@ async function createKernelSession(
restartCount: 0, restartCount: 0,
dead: false, dead: false,
needsRestart: false, needsRestart: false,
kernelInvalidatedByRecovery: false,
disposing: false, disposing: false,
disposeResultPromise: undefined, disposeResultPromise: undefined,
nextDisposalRetryAt: undefined, nextDisposalRetryAt: undefined,
@@ -621,18 +551,16 @@ async function restartKernelSession(
} }
requireRemainingTimeoutMs(options.deadlineMs); requireRemainingTimeoutMs(options.deadlineMs);
try { try {
if (!session.kernelInvalidatedByRecovery) { const deadKernel = session.dead || !session.kernel.isAlive();
const deadKernel = session.dead || !session.kernel.isAlive(); const shutdownTimeoutMs = requireRemainingTimeoutMs(options.deadlineMs);
const shutdownTimeoutMs = requireRemainingTimeoutMs(options.deadlineMs); const shutdownResult = await session.kernel.shutdown({ signal: options.signal, timeoutMs: shutdownTimeoutMs });
const shutdownResult = await session.kernel.shutdown({ signal: options.signal, timeoutMs: shutdownTimeoutMs }); if (!shutdownResult.confirmed && !deadKernel) {
if (!shutdownResult.confirmed && !deadKernel) { throw new Error("Failed to confirm crashed kernel shutdown before restart");
throw new Error("Failed to confirm crashed kernel shutdown before restart"); }
} if (!shutdownResult.confirmed) {
if (!shutdownResult.confirmed) { logger.warn("Proceeding with retained kernel restart after unconfirmed dead-kernel shutdown", {
logger.warn("Proceeding with retained kernel restart after unconfirmed dead-kernel shutdown", { sessionId: session.id,
sessionId: session.id, });
});
}
} }
const env = buildKernelEnv(options); const env = buildKernelEnv(options);
const startOptions = buildKernelStartOptions(cwd, env, options); const startOptions = buildKernelStartOptions(cwd, env, options);
@@ -640,7 +568,6 @@ async function restartKernelSession(
session.kernel = kernel; session.kernel = kernel;
session.dead = false; session.dead = false;
session.needsRestart = false; session.needsRestart = false;
session.kernelInvalidatedByRecovery = false;
session.lastUsedAt = Date.now(); session.lastUsedAt = Date.now();
ensureKernelHeartbeat(session); ensureKernelHeartbeat(session);
} catch (err) { } catch (err) {
@@ -654,9 +581,6 @@ type KernelDisposalResult = { status: "confirmed" } | { status: "unconfirmed" }
type KernelDisposalWaitResult = KernelDisposalResult | { status: "timedOut" }; type KernelDisposalWaitResult = KernelDisposalResult | { status: "timedOut" };
function createKernelDisposalResultPromise(session: KernelSession, timeoutMs?: number): Promise<KernelDisposalResult> { function createKernelDisposalResultPromise(session: KernelSession, timeoutMs?: number): Promise<KernelDisposalResult> {
if (session.kernelInvalidatedByRecovery) {
return Promise.resolve({ status: "confirmed" as const });
}
return Promise.resolve() return Promise.resolve()
.then(() => session.kernel.shutdown(timeoutMs === undefined ? undefined : { timeoutMs })) .then(() => session.kernel.shutdown(timeoutMs === undefined ? undefined : { timeoutMs }))
.then( .then(
@@ -1,424 +0,0 @@
import * as fs from "node:fs";
import { createServer } from "node:net";
import * as path from "node:path";
import { Process } from "@oh-my-pi/pi-natives";
import { getPythonGatewayDir, isEnoent, logger, procmgr } from "@oh-my-pi/pi-utils";
import type { Subprocess } from "bun";
import { Settings } from "../../config/settings";
import { getOrCreateSnapshot } from "../../utils/shell-snapshot";
import { filterEnv, resolvePythonRuntime } from "./runtime";
const GATEWAY_INFO_FILE = "gateway.json";
const GATEWAY_LOCK_FILE = "gateway.lock";
const GATEWAY_STARTUP_TIMEOUT_MS = 30000;
const GATEWAY_LOCK_TIMEOUT_MS = GATEWAY_STARTUP_TIMEOUT_MS + 5000;
const GATEWAY_LOCK_RETRY_MS = 50;
const GATEWAY_LOCK_STALE_MS = GATEWAY_STARTUP_TIMEOUT_MS * 2;
const GATEWAY_LOCK_HEARTBEAT_MS = 5000;
const HEALTH_CHECK_TIMEOUT_MS = 3000;
export interface GatewayInfo {
url: string;
pid: number;
startedAt: number;
pythonPath?: string;
venvPath?: string | null;
}
interface GatewayLockInfo {
pid: number;
startedAt: number;
}
interface AcquireResult {
url: string;
isShared: boolean;
}
let localGatewayProcess: Subprocess | null = null;
let localGatewayUrl: string | null = null;
let isCoordinatorInitialized = false;
async function allocatePort(): Promise<number> {
const { promise, resolve, reject } = Promise.withResolvers<number>();
const server = createServer();
server.unref();
server.on("error", reject);
server.listen(0, "127.0.0.1", () => {
const address = server.address();
if (address && typeof address === "object") {
const port = address.port;
server.close((err: Error | null | undefined) => {
if (err) {
reject(err);
} else {
resolve(port);
}
});
} else {
server.close();
reject(new Error("Failed to allocate port"));
}
});
return promise;
}
function getGatewayDir(): string {
return getPythonGatewayDir();
}
function getGatewayInfoPath(): string {
return path.join(getGatewayDir(), GATEWAY_INFO_FILE);
}
function getGatewayLockPath(): string {
return path.join(getGatewayDir(), GATEWAY_LOCK_FILE);
}
async function writeLockInfo(lockPath: string): Promise<void> {
const payload: GatewayLockInfo = { pid: process.pid, startedAt: Date.now() };
try {
await Bun.write(lockPath, JSON.stringify(payload));
} catch {
// Ignore lock write failures
}
}
async function readLockInfo(lockPath: string): Promise<GatewayLockInfo | null> {
try {
const raw = await Bun.file(lockPath).text();
const parsed = JSON.parse(raw) as Partial<GatewayLockInfo>;
if (typeof parsed.pid === "number" && Number.isFinite(parsed.pid)) {
return { pid: parsed.pid, startedAt: typeof parsed.startedAt === "number" ? parsed.startedAt : 0 };
}
} catch {
// Ignore parse errors
}
return null;
}
async function ensureGatewayDir(): Promise<void> {
const dir = getGatewayDir();
await fs.promises.mkdir(dir, { recursive: true });
}
async function withGatewayLock<T>(handler: () => Promise<T>): Promise<T> {
await ensureGatewayDir();
const lockPath = getGatewayLockPath();
const start = Date.now();
while (true) {
let fd: fs.promises.FileHandle | undefined;
try {
fd = await fs.promises.open(lockPath, "wx");
let heartbeatRunning = true;
const heartbeat = (async () => {
while (heartbeatRunning) {
await Bun.sleep(GATEWAY_LOCK_HEARTBEAT_MS);
if (!heartbeatRunning) break;
try {
const now = new Date();
await fs.promises.utimes(lockPath, now, now);
} catch {
// Ignore heartbeat errors
}
}
})();
try {
await writeLockInfo(lockPath);
return await handler();
} finally {
heartbeatRunning = false;
void heartbeat.catch(() => {}); // Don't await - let it die naturally
try {
await fd.close();
await fs.promises.unlink(lockPath);
} catch {
// Ignore lock cleanup errors
}
}
} catch (err) {
const error = err as NodeJS.ErrnoException;
if (error.code === "EEXIST") {
let removedStale = false;
try {
const lockStat = await fs.promises.stat(lockPath);
const lockInfo = await readLockInfo(lockPath);
const lockPid = lockInfo?.pid;
const lockAgeMs = lockInfo?.startedAt ? Date.now() - lockInfo.startedAt : Date.now() - lockStat.mtimeMs;
const staleByTime = lockAgeMs > GATEWAY_LOCK_STALE_MS;
const staleByPid = lockPid !== undefined && !procmgr.isPidRunning(lockPid);
const staleByMissingPid = lockPid === undefined && staleByTime;
if (staleByPid || staleByMissingPid) {
await fs.promises.unlink(lockPath);
removedStale = true;
logger.warn("Removed stale shared gateway lock", { path: lockPath, pid: lockPid });
}
} catch {
// Ignore stat errors; keep waiting
}
if (!removedStale) {
if (Date.now() - start > GATEWAY_LOCK_TIMEOUT_MS) {
throw new Error("Timed out waiting for shared gateway lock");
}
await Bun.sleep(GATEWAY_LOCK_RETRY_MS);
}
continue;
}
throw err;
}
}
}
async function readGatewayInfo(): Promise<GatewayInfo | null> {
const infoPath = getGatewayInfoPath();
try {
const content = await Bun.file(infoPath).text();
const parsed = JSON.parse(content) as Partial<GatewayInfo>;
if (typeof parsed.url !== "string" || typeof parsed.pid !== "number" || typeof parsed.startedAt !== "number") {
return null;
}
return {
url: parsed.url,
pid: parsed.pid,
startedAt: parsed.startedAt,
pythonPath: typeof parsed.pythonPath === "string" ? parsed.pythonPath : undefined,
venvPath: typeof parsed.venvPath === "string" || parsed.venvPath === null ? parsed.venvPath : undefined,
};
} catch (err) {
if (isEnoent(err)) return null;
return null;
}
}
async function writeGatewayInfo(info: GatewayInfo): Promise<void> {
const infoPath = getGatewayInfoPath();
const tempPath = `${infoPath}.tmp`;
await Bun.write(tempPath, JSON.stringify(info, null, 2));
await fs.promises.rename(tempPath, infoPath);
}
async function clearGatewayInfo(): Promise<void> {
const infoPath = getGatewayInfoPath();
try {
await fs.promises.unlink(infoPath);
} catch {
// Ignore errors on cleanup (file may not exist)
}
}
async function isGatewayHealthy(url: string): Promise<boolean> {
try {
const response = await fetch(`${url}/api/kernelspecs`, {
signal: AbortSignal.timeout(HEALTH_CHECK_TIMEOUT_MS),
});
return response.ok;
} catch {
return false;
}
}
async function isGatewayAlive(info: GatewayInfo): Promise<boolean> {
if (!procmgr.isPidRunning(info.pid)) return false;
return await isGatewayHealthy(info.url);
}
async function startGatewayProcess(
cwd: string,
): Promise<{ url: string; pid: number; pythonPath: string; venvPath: string | null }> {
const settings = await Settings.init();
const { shell, env } = settings.getShellConfig();
const filteredEnv = filterEnv(env);
const runtime = resolvePythonRuntime(cwd, filteredEnv);
const snapshotPath = await getOrCreateSnapshot(shell, env).catch((err: unknown) => {
logger.warn("Failed to resolve shell snapshot for shared Python gateway", {
error: err instanceof Error ? err.message : String(err),
});
return null;
});
const kernelEnv: Record<string, string | undefined> = {
...runtime.env,
PYTHONUNBUFFERED: "1",
PI_SHELL_SNAPSHOT: snapshotPath ?? undefined,
};
const gatewayPort = await allocatePort();
const gatewayUrl = `http://127.0.0.1:${gatewayPort}`;
const gatewayProcess = Bun.spawn(
[
runtime.pythonPath,
"-m",
"kernel_gateway",
"--KernelGatewayApp.ip=127.0.0.1",
`--KernelGatewayApp.port=${gatewayPort}`,
"--KernelGatewayApp.port_retries=0",
"--KernelGatewayApp.allow_origin=*",
"--JupyterApp.answer_yes=true",
],
{
cwd,
stdin: "ignore",
stdout: "pipe",
stderr: "pipe",
windowsHide: true,
detached: true,
env: kernelEnv,
},
);
let exited = false;
gatewayProcess.exited
.catch(() => {})
.then(() => {
exited = true;
});
const startTime = Date.now();
while (Date.now() - startTime < GATEWAY_STARTUP_TIMEOUT_MS) {
if (exited) {
throw new Error("Gateway process exited during startup");
}
if (await isGatewayHealthy(gatewayUrl)) {
localGatewayProcess = gatewayProcess;
localGatewayUrl = gatewayUrl;
return {
url: gatewayUrl,
pid: gatewayProcess.pid,
pythonPath: runtime.pythonPath,
venvPath: runtime.venvPath ?? null,
};
}
await Bun.sleep(100);
}
gatewayProcess.kill();
throw new Error("Gateway startup timeout");
}
async function killGateway(pid: number, context: string): Promise<void> {
try {
await Process.fromPid(pid)?.terminate();
} catch (err) {
logger.warn("Failed to kill shared gateway process", {
error: err instanceof Error ? err.message : String(err),
pid,
context,
});
}
}
export async function acquireSharedGateway(cwd: string): Promise<AcquireResult | null> {
try {
return await withGatewayLock(async () => {
const existingInfo = await logger.time("acquireSharedGateway:readInfo", readGatewayInfo);
if (existingInfo) {
if (await logger.time("acquireSharedGateway:isAlive", isGatewayAlive, existingInfo)) {
localGatewayUrl = existingInfo.url;
isCoordinatorInitialized = true;
logger.debug("Reusing global Python gateway", { url: existingInfo.url });
return { url: existingInfo.url, isShared: true };
}
logger.debug("Cleaning up stale gateway info", { pid: existingInfo.pid });
if (procmgr.isPidRunning(existingInfo.pid)) {
await killGateway(existingInfo.pid, "stale");
}
await clearGatewayInfo();
}
const { url, pid, pythonPath, venvPath } = await logger.time(
"acquireSharedGateway:startGateway",
startGatewayProcess,
cwd,
);
const info: GatewayInfo = {
url,
pid,
startedAt: Date.now(),
pythonPath,
venvPath,
};
await writeGatewayInfo(info);
isCoordinatorInitialized = true;
logger.debug("Started global Python gateway", { url, pid });
return { url, isShared: true };
});
} catch (err) {
logger.warn("Failed to acquire shared gateway, falling back to local", {
error: err instanceof Error ? err.message : String(err),
});
return null;
}
}
export async function releaseSharedGateway(): Promise<void> {
if (!isCoordinatorInitialized) return;
}
export async function getSharedGatewayUrl(): Promise<string | null> {
if (localGatewayUrl) return localGatewayUrl;
return (await readGatewayInfo())?.url ?? null;
}
export async function isSharedGatewayActive(): Promise<boolean> {
return (await getGatewayStatus()).active;
}
export interface GatewayStatus {
active: boolean;
url: string | null;
pid: number | null;
uptime: number | null;
pythonPath: string | null;
venvPath: string | null;
}
export async function getGatewayStatus(): Promise<GatewayStatus> {
const info = await readGatewayInfo();
if (!info) {
return {
active: false,
url: null,
pid: null,
uptime: null,
pythonPath: null,
venvPath: null,
};
}
const active = procmgr.isPidRunning(info.pid);
return {
active,
url: info.url,
pid: info.pid,
uptime: active ? Date.now() - info.startedAt : null,
pythonPath: info.pythonPath ?? null,
venvPath: info.venvPath ?? null,
};
}
export async function shutdownSharedGateway(): Promise<void> {
try {
await withGatewayLock(async () => {
const info = await readGatewayInfo();
if (!info) return;
if (procmgr.isPidRunning(info.pid)) {
await killGateway(info.pid, "shutdown");
}
await clearGatewayInfo();
});
} catch (err) {
logger.warn("Failed to shutdown shared gateway", {
error: err instanceof Error ? err.message : String(err),
});
} finally {
if (localGatewayProcess) {
await killGateway(localGatewayProcess.pid, "shutdown-local");
}
localGatewayProcess = null;
localGatewayUrl = null;
isCoordinatorInitialized = false;
}
}
@@ -25,7 +25,6 @@ export default {
}, },
async execute(code: string, opts: ExecutorBackendExecOptions): Promise<ExecutorBackendResult> { async execute(code: string, opts: ExecutorBackendExecOptions): Promise<ExecutorBackendResult> {
const useSharedGateway = readSetting<boolean>(opts.session, "python.sharedGateway");
const kernelMode = readSetting<PythonExecutorOptions["kernelMode"]>(opts.session, "python.kernelMode"); const kernelMode = readSetting<PythonExecutorOptions["kernelMode"]>(opts.session, "python.kernelMode");
const executorOptions: PythonExecutorOptions = { const executorOptions: PythonExecutorOptions = {
cwd: opts.cwd, cwd: opts.cwd,
@@ -33,7 +32,6 @@ export default {
signal: opts.signal, signal: opts.signal,
sessionId: namespaceSessionId(opts.sessionId), sessionId: namespaceSessionId(opts.sessionId),
kernelMode, kernelMode,
useSharedGateway,
sessionFile: opts.sessionFile, sessionFile: opts.sessionFile,
artifactsDir: opts.session.getArtifactsDir?.() ?? undefined, artifactsDir: opts.session.getArtifactsDir?.() ?? undefined,
kernelOwnerId: opts.kernelOwnerId, kernelOwnerId: opts.kernelOwnerId,
File diff suppressed because it is too large Load Diff
+11 -7
View File
@@ -1,10 +1,13 @@
from __future__ import annotations from __future__ import annotations
# OMP IPython prelude helpers # OMP prelude helpers (loaded once into the runner namespace)
if "__omp_prelude_loaded__" not in globals(): if "__omp_prelude_loaded__" not in globals():
__omp_prelude_loaded__ = True __omp_prelude_loaded__ = True
from pathlib import Path from pathlib import Path
import os, json import os, json
from IPython.display import display as _ipy_display, JSON
# __omp_display is injected by runner.py before the prelude executes; it
# mirrors IPython's display() semantics with the same MIME bundle output.
_omp_display = __omp_display # type: ignore[name-defined]
_PRESENTABLE_REPRS = ( _PRESENTABLE_REPRS = (
"_repr_mimebundle_", "_repr_mimebundle_",
@@ -18,21 +21,22 @@ if "__omp_prelude_loaded__" not in globals():
) )
def display(value): def display(value):
"""Render a value. Wraps plain dict/list values as interactive JSON.""" """Render a value. Falls back to a JSON+text/plain bundle for plain dict/list/tuple."""
if any(hasattr(value, attr) for attr in _PRESENTABLE_REPRS): if any(hasattr(value, attr) for attr in _PRESENTABLE_REPRS):
_ipy_display(value) _omp_display(value)
return return
if isinstance(value, (dict, list, tuple)): if isinstance(value, (dict, list, tuple)):
try: try:
_ipy_display(JSON(value)) bundle = {"application/json": value, "text/plain": repr(value)}
_omp_display(bundle, raw=True)
return return
except Exception: except Exception:
pass pass
_ipy_display(value) _omp_display(value)
def _emit_status(op: str, **data): def _emit_status(op: str, **data):
"""Emit structured status event for TUI rendering.""" """Emit structured status event for TUI rendering."""
_ipy_display({"application/x-omp-status": {"op": op, **data}}, raw=True) _omp_display({"application/x-omp-status": {"op": op, **data}}, raw=True)
def env(key: str | None = None, value: str | None = None): def env(key: str | None = None, value: str | None = None):
+879
View File
@@ -0,0 +1,879 @@
"""OMP Python runner — subprocess wrapper used by the coding-agent host.
NDJSON protocol over stdin/stdout. Host writes one JSON object per line;
wrapper writes typed frames back.
Host -> wrapper:
{"id": str, "code": str, "silent": bool?, "storeHistory": bool?}
{"type": "exit"} # graceful shutdown
Wrapper -> host:
{"type": "started", "id": ...}
{"type": "stdout", "id": ..., "data": str}
{"type": "stderr", "id": ..., "data": str}
{"type": "display", "id": ..., "bundle": {<mime>: <value>}}
{"type": "result", "id": ..., "bundle": {<mime>: <value>}}
{"type": "error", "id": ..., "ename": str, "evalue": str, "traceback": [str]}
{"type": "done", "id": ..., "status": "ok"|"error",
"executionCount": int, "cancelled": bool}
The runner is intentionally self-contained: no third-party imports, no IPython.
Magics are translated by a small line-scanner before AST parsing; rich display
falls back through `_repr_*_` methods so pandas/PIL/plotly etc. still render
when installed.
"""
from __future__ import annotations
import ast
import base64
import builtins
import io
import json
import os
import re
import runpy
import shlex
import signal
import subprocess
import sys
import threading
import time
import traceback
from pathlib import Path
from typing import Any, Callable
# ---------------------------------------------------------------------------
# Frame writer
# ---------------------------------------------------------------------------
_RAW_STDOUT = sys.__stdout__
_RAW_STDERR = sys.__stderr__
_OUT_LOCK = threading.Lock()
def _json_default(o: Any) -> Any:
try:
return repr(o)
except Exception:
return f"<unrepr {type(o).__name__}>"
def _emit(frame: dict) -> None:
"""Serialize a frame and write it to the host as a single NDJSON line."""
line = json.dumps(frame, ensure_ascii=False, default=_json_default)
with _OUT_LOCK:
_RAW_STDOUT.write(line)
_RAW_STDOUT.write("\n")
_RAW_STDOUT.flush()
# ---------------------------------------------------------------------------
# User stdout/stderr proxies
# ---------------------------------------------------------------------------
class _StreamProxy(io.TextIOBase):
"""Emit each ``write()`` as a typed frame tied to the current request."""
def __init__(self, kind: str) -> None:
super().__init__()
self._kind = kind
def writable(self) -> bool: # noqa: D401 - protocol method
return True
def isatty(self) -> bool: # noqa: D401 - protocol method
return False
def write(self, data: Any) -> int: # type: ignore[override]
if not isinstance(data, str):
data = str(data)
if not data:
return 0
rid = _STATE.current_id
if rid is None:
_RAW_STDERR.write(data)
_RAW_STDERR.flush()
return len(data)
_emit({"type": self._kind, "id": rid, "data": data})
return len(data)
def flush(self) -> None: # noqa: D401 - protocol method
return None
# ---------------------------------------------------------------------------
# Runner state
# ---------------------------------------------------------------------------
class _RunnerState:
def __init__(self) -> None:
self.current_id: str | None = None
self.execution_count: int = 0
self.cancel_requested: bool = False
# User globals — kept across requests when running in session mode.
self.user_ns: dict[str, Any] = {
"__name__": "__main__",
"__doc__": None,
"__builtins__": builtins,
}
self.last_install_marker: int = 0
_STATE = _RunnerState()
# ---------------------------------------------------------------------------
# Magic source transformer
# ---------------------------------------------------------------------------
_MAGIC_LINE_RE = re.compile(r"^(?P<indent>[ \t]*)(?P<name>[A-Za-z_][A-Za-z_0-9]*)(?:[ \t]+(?P<args>.*))?$")
_ASSIGN_LINE_RE = re.compile(
r"^(?P<indent>[ \t]*)(?P<lhs>[A-Za-z_][A-Za-z_0-9.\[\], ]*?)\s*=\s*(?P<rhs>.+)$"
)
def _fold_continuations(lines: list[str], start: int) -> tuple[str, int]:
"""Fold trailing backslash continuations starting at ``start``. Returns
``(folded_text, lines_consumed)``."""
parts: list[str] = []
i = start
while i < len(lines):
line = lines[i]
if line.endswith("\\"):
parts.append(line[:-1])
i += 1
continue
parts.append(line)
i += 1
break
return ("".join(parts), i - start)
def _quote_arg(text: str) -> str:
"""Return a Python string literal that round-trips ``text`` exactly."""
return json.dumps(text, ensure_ascii=False)
def transform_cell(source: str) -> str:
"""Translate IPython-style magics + shell escapes into plain Python.
Rules
-----
* ``%name args`` -> ``__omp_magic("name", "args")``
* ``var = %name args`` -> ``var = __omp_magic("name", "args")``
* ``!cmd`` -> ``__omp_shell("cmd")``
* ``var = !cmd`` -> ``var = __omp_shell("cmd")``
* ``%%name args\\n<body>`` -> ``__omp_magic_cell("name", "args", "<body>")``
(cell magic must be the first non-whitespace token of a top-level line and
consumes the remainder of the cell)
Lines inside strings or comments are left alone — we operate on the raw
text before parsing, but the scanner only fires on the first token of each
physical line and never touches the body of triple-quoted strings because
those bodies are never first tokens themselves.
"""
if "%" not in source and "!" not in source:
return source
lines = source.splitlines()
out: list[str] = []
i = 0
while i < len(lines):
line = lines[i]
stripped = line.lstrip()
indent = line[: len(line) - len(stripped)]
# Cell magic — consumes from here to EOF.
if stripped.startswith("%%"):
head, _ = _split_magic_head(stripped[2:])
name, args = head
body_lines = lines[i + 1 :]
body = "\n".join(body_lines)
out.append(
f"{indent}__omp_magic_cell({_quote_arg(name)}, {_quote_arg(args)}, {_quote_arg(body)})"
)
return "\n".join(out)
# Line magic / shell at start of line.
if stripped.startswith("%") and not stripped.startswith("%%"):
folded, consumed = _fold_continuations(lines, i)
stripped_folded = folded.lstrip()
indent = folded[: len(folded) - len(stripped_folded)]
head, _ = _split_magic_head(stripped_folded[1:])
name, args = head
out.append(f"{indent}__omp_magic({_quote_arg(name)}, {_quote_arg(args)})")
i += consumed
continue
if stripped.startswith("!"):
folded, consumed = _fold_continuations(lines, i)
stripped_folded = folded.lstrip()
indent = folded[: len(folded) - len(stripped_folded)]
cmd = stripped_folded[1:].strip()
out.append(f"{indent}__omp_shell({_quote_arg(cmd)})")
i += consumed
continue
# Assignment forms: var = %magic / var = !cmd
m = _ASSIGN_LINE_RE.match(line)
if m:
rhs = m.group("rhs").strip()
if rhs.startswith("!"):
cmd = rhs[1:].strip()
out.append(f"{m.group('indent')}{m.group('lhs').rstrip()} = __omp_shell({_quote_arg(cmd)})")
i += 1
continue
if rhs.startswith("%") and not rhs.startswith("%%"):
head, _ = _split_magic_head(rhs[1:])
name, args = head
out.append(
f"{m.group('indent')}{m.group('lhs').rstrip()} = __omp_magic({_quote_arg(name)}, {_quote_arg(args)})"
)
i += 1
continue
out.append(line)
i += 1
return "\n".join(out)
def _split_magic_head(text: str) -> tuple[tuple[str, str], str]:
"""Split ``"name rest"`` into ``("name", "rest")``."""
text = text.lstrip()
if not text:
return ("", ""), ""
m = re.match(r"([A-Za-z_][A-Za-z_0-9]*)(?:\s+(.*))?$", text)
if not m:
return ("", text), ""
return (m.group(1), (m.group(2) or "").rstrip()), ""
# ---------------------------------------------------------------------------
# Magic registry
# ---------------------------------------------------------------------------
_LINE_MAGICS: dict[str, Callable[[str], Any]] = {}
_CELL_MAGICS: dict[str, Callable[[str, str], Any]] = {}
def line_magic(name: str) -> Callable[[Callable[[str], Any]], Callable[[str], Any]]:
def decorator(fn: Callable[[str], Any]) -> Callable[[str], Any]:
_LINE_MAGICS[name] = fn
return fn
return decorator
def cell_magic(name: str) -> Callable[[Callable[[str, str], Any]], Callable[[str, str], Any]]:
def decorator(fn: Callable[[str, str], Any]) -> Callable[[str, str], Any]:
_CELL_MAGICS[name] = fn
return fn
return decorator
def _emit_status(op: str, **data: Any) -> None:
bundle = {"application/x-omp-status": {"op": op, **data}}
rid = _STATE.current_id
if rid is None:
return
_emit({"type": "display", "id": rid, "bundle": bundle})
@line_magic("pip")
def _magic_pip(args: str) -> None:
argv = shlex.split(args) if args else ["--help"]
cmd = [sys.executable, "-m", "pip", *argv]
proc = subprocess.Popen(
cmd,
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
text=True,
bufsize=1,
)
installed_packages: list[str] = []
assert proc.stdout is not None
for raw_line in proc.stdout:
sys.stdout.write(raw_line)
m = re.search(r"Successfully installed\s+(.+)$", raw_line)
if m:
for token in m.group(1).split():
# Token is name-version; drop the version suffix.
pkg = token.rsplit("-", 1)[0]
installed_packages.append(pkg.replace("_", "-"))
proc.wait()
if installed_packages:
import importlib
importlib.invalidate_caches()
prefixes = {pkg.lower().replace("-", "_") for pkg in installed_packages}
for mod_name in list(sys.modules):
head = mod_name.split(".", 1)[0].lower()
if head in prefixes:
sys.modules.pop(mod_name, None)
_emit_status("pip", args=args, installed=installed_packages, exit_code=proc.returncode)
@line_magic("cd")
def _magic_cd(args: str) -> str:
path = os.path.expanduser(args.strip()) or os.path.expanduser("~")
os.chdir(path)
cwd = os.getcwd()
_emit_status("cd", path=cwd)
return cwd
@line_magic("pwd")
def _magic_pwd(_args: str) -> str:
cwd = os.getcwd()
_emit_status("pwd", path=cwd)
return cwd
@line_magic("ls")
def _magic_ls(args: str) -> list[str]:
target = os.path.expanduser(args.strip()) or "."
entries = sorted(os.listdir(target))
_emit_status("ls", path=os.path.abspath(target), count=len(entries))
return entries
@line_magic("env")
def _magic_env(args: str) -> Any:
args = args.strip()
if not args:
return dict(sorted(os.environ.items()))
if "=" in args:
key, value = args.split("=", 1)
os.environ[key.strip()] = value.strip()
return value.strip()
return os.environ.get(args)
@line_magic("set_env")
def _magic_set_env(args: str) -> str:
parts = args.split(None, 1)
if len(parts) != 2:
raise ValueError("Usage: %set_env KEY VALUE")
key, value = parts
os.environ[key] = value
return value
@line_magic("time")
def _magic_time(args: str) -> Any:
start = time.perf_counter()
result = eval(args, _STATE.user_ns)
elapsed = time.perf_counter() - start
sys.stdout.write(f"Wall time: {elapsed * 1000:.2f} ms\n")
_emit_status("time", elapsed_ms=round(elapsed * 1000, 3))
return result
@line_magic("timeit")
def _magic_timeit(args: str) -> None:
import timeit as _timeit
timer = _timeit.Timer(stmt=args, globals=_STATE.user_ns)
iters, total = timer.autorange()
per = total / iters
sys.stdout.write(f"{iters} loops, best of 1: {per * 1e6:.2f} us per loop\n")
_emit_status("timeit", loops=iters, total_ms=round(total * 1000, 3))
@line_magic("who")
def _magic_who(_args: str) -> list[str]:
names = sorted(
name
for name, value in _STATE.user_ns.items()
if not name.startswith("_") and not callable(value) or hasattr(value, "__class__")
)
return [n for n in names if not n.startswith("__")]
@line_magic("whos")
def _magic_whos(_args: str) -> list[tuple[str, str]]:
rows = []
for name in sorted(_STATE.user_ns):
if name.startswith("__"):
continue
value = _STATE.user_ns[name]
rows.append((name, type(value).__name__))
return rows
@line_magic("reset")
def _magic_reset(_args: str) -> None:
_STATE.user_ns.clear()
_STATE.user_ns.update({"__name__": "__main__", "__doc__": None, "__builtins__": builtins})
_install_builtins(_STATE.user_ns)
_emit_status("reset")
@line_magic("load")
def _magic_load(args: str) -> None:
path = Path(os.path.expanduser(args.strip()))
source = path.read_text(encoding="utf-8")
_emit({"type": "display", "id": _STATE.current_id, "bundle": {"text/plain": source}})
_exec_source(source, _STATE.user_ns)
@line_magic("run")
def _magic_run(args: str) -> None:
parts = shlex.split(args) if args else []
if not parts:
raise ValueError("Usage: %run <path>")
target = os.path.expanduser(parts[0])
saved_argv = sys.argv
try:
sys.argv = [target, *parts[1:]]
result_ns = runpy.run_path(target, run_name="__main__")
finally:
sys.argv = saved_argv
for name, value in result_ns.items():
if name.startswith("__"):
continue
_STATE.user_ns[name] = value
@cell_magic("bash")
def _magic_cell_bash(args: str, body: str) -> int:
return _run_shell_body(body, shell_arg="/bin/bash")
@cell_magic("sh")
def _magic_cell_sh(args: str, body: str) -> int:
return _run_shell_body(body, shell_arg="/bin/sh")
@cell_magic("capture")
def _magic_cell_capture(args: str, body: str) -> str:
"""Capture stdout/stderr of body; bind to ``args`` (a name) if provided."""
captured = io.StringIO()
saved_stdout, saved_stderr = sys.stdout, sys.stderr
sys.stdout = sys.stderr = captured
try:
_exec_source(body, _STATE.user_ns)
finally:
sys.stdout, sys.stderr = saved_stdout, saved_stderr
text = captured.getvalue()
name = args.strip()
if name:
_STATE.user_ns[name] = text
return text
@cell_magic("timeit")
def _magic_cell_timeit(args: str, body: str) -> None:
import timeit as _timeit
timer = _timeit.Timer(stmt=body, globals=_STATE.user_ns)
iters, total = timer.autorange()
per = total / iters
sys.stdout.write(f"{iters} loops, best of 1: {per * 1e6:.2f} us per loop\n")
_emit_status("timeit", loops=iters, total_ms=round(total * 1000, 3))
@cell_magic("writefile")
def _magic_cell_writefile(args: str, body: str) -> str:
path = Path(os.path.expanduser(args.strip()))
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(body, encoding="utf-8")
_emit_status("writefile", path=str(path), bytes=len(body))
return str(path)
def _run_shell_body(body: str, *, shell_arg: str) -> int:
proc = subprocess.Popen(
[shell_arg, "-c", body],
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
text=True,
bufsize=1,
)
assert proc.stdout is not None
for raw_line in proc.stdout:
sys.stdout.write(raw_line)
proc.wait()
return proc.returncode
def __omp_magic(name: str, args: str) -> Any:
fn = _LINE_MAGICS.get(name)
if fn is None:
raise NameError(f"UsageError: Line magic function '%{name}' not found.")
return fn(args)
def __omp_magic_cell(name: str, args: str, body: str) -> Any:
fn = _CELL_MAGICS.get(name)
if fn is None:
raise NameError(f"UsageError: Cell magic function '%%{name}' not found.")
return fn(args, body)
class _ShellResult(list):
"""Result of ``!cmd`` — list of stripped output lines."""
def __init__(self, lines: list[str], returncode: int) -> None:
super().__init__(lines)
self.returncode = returncode
@property
def n(self) -> str: # IPython compat
return "\n".join(self)
@property
def s(self) -> str: # IPython compat
return " ".join(self)
def __omp_shell(cmd: str) -> _ShellResult:
proc = subprocess.run(
cmd,
shell=True,
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
text=True,
)
if proc.stdout:
sys.stdout.write(proc.stdout)
lines = [line for line in (proc.stdout or "").splitlines()]
return _ShellResult(lines, proc.returncode)
# ---------------------------------------------------------------------------
# Display dispatch
# ---------------------------------------------------------------------------
_REPR_MIMES = [
("_repr_html_", "text/html"),
("_repr_markdown_", "text/markdown"),
("_repr_svg_", "image/svg+xml"),
("_repr_png_", "image/png"),
("_repr_jpeg_", "image/jpeg"),
("_repr_json_", "application/json"),
("_repr_latex_", "text/latex"),
]
def _coerce_image_bytes(value: Any) -> str:
if isinstance(value, (bytes, bytearray)):
return base64.b64encode(bytes(value)).decode("ascii")
if isinstance(value, str):
return value
return base64.b64encode(repr(value).encode("utf-8")).decode("ascii")
def _mime_bundle(value: Any) -> dict:
"""Build a Jupyter-style MIME bundle for ``value``.
Honors ``_repr_mimebundle_`` first, falls back to individual ``_repr_*_``
accessors, and always provides ``text/plain``.
"""
bundle: dict[str, Any] = {}
mimebundle = getattr(value, "_repr_mimebundle_", None)
if callable(mimebundle):
try:
data = mimebundle()
except Exception:
data = None
if isinstance(data, tuple):
data = data[0]
if isinstance(data, dict):
bundle.update({str(k): v for k, v in data.items()})
for attr, mime in _REPR_MIMES:
if mime in bundle:
continue
repr_fn = getattr(value, attr, None)
if not callable(repr_fn):
continue
try:
data = repr_fn()
except Exception:
continue
if data is None:
continue
if mime in ("image/png", "image/jpeg"):
bundle[mime] = _coerce_image_bytes(data)
else:
bundle[mime] = data
if "text/plain" not in bundle:
try:
bundle["text/plain"] = repr(value)
except Exception:
bundle["text/plain"] = f"<unrepr {type(value).__name__}>"
return bundle
def _emit_display(bundle: dict, *, kind: str = "display") -> None:
rid = _STATE.current_id
if rid is None:
return
_emit({"type": kind, "id": rid, "bundle": bundle})
def __omp_display(value: Any, *, raw: bool = False, kind: str = "display") -> None:
if raw:
if not isinstance(value, dict):
raise TypeError("display(..., raw=True) requires a MIME bundle dict")
bundle = {str(k): v for k, v in value.items()}
if "text/plain" not in bundle:
bundle["text/plain"] = ""
_emit_display(bundle, kind=kind)
return
_emit_display(_mime_bundle(value), kind=kind)
# ---------------------------------------------------------------------------
# Matplotlib post-cell flush
# ---------------------------------------------------------------------------
def _flush_matplotlib_figures() -> None:
plt = sys.modules.get("matplotlib.pyplot")
if plt is None:
return
try:
fignums = list(plt.get_fignums())
except Exception:
return
for num in fignums:
try:
fig = plt.figure(num)
buf = io.BytesIO()
fig.savefig(buf, format="png", bbox_inches="tight")
data = base64.b64encode(buf.getvalue()).decode("ascii")
_emit_display({"image/png": data, "text/plain": f"<Figure {num}>"})
plt.close(fig)
except Exception:
continue
# Force a non-interactive backend before user code imports matplotlib. Set as
# environ default so the user can still override it explicitly.
os.environ.setdefault("MPLBACKEND", "Agg")
# ---------------------------------------------------------------------------
# Builtin injection
# ---------------------------------------------------------------------------
def _install_builtins(ns: dict) -> None:
ns["display"] = __omp_display
ns["__omp_display"] = __omp_display
ns["__omp_magic"] = __omp_magic
ns["__omp_magic_cell"] = __omp_magic_cell
ns["__omp_shell"] = __omp_shell
_install_builtins(_STATE.user_ns)
# ---------------------------------------------------------------------------
# Source execution (split last expression for rich display)
# ---------------------------------------------------------------------------
def _exec_source(source: str, ns: dict) -> None:
"""Compile + execute ``source``; if the last node is an expression, route
its value through ``__omp_display`` so dataframes/figures render rich."""
try:
module = ast.parse(source, mode="exec")
except SyntaxError:
raise
if not module.body:
return
last = module.body[-1]
if isinstance(last, ast.Expr):
body_module = ast.Module(body=module.body[:-1], type_ignores=[])
expr_module = ast.Expression(body=last.value)
ast.copy_location(expr_module, last)
body_code = compile(body_module, "<cell>", "exec")
expr_code = compile(expr_module, "<cell>", "eval")
exec(body_code, ns)
value = eval(expr_code, ns)
if value is not None:
__omp_display(value, kind="result")
return
code = compile(module, "<cell>", "exec")
exec(code, ns)
# ---------------------------------------------------------------------------
# Signal handling
# ---------------------------------------------------------------------------
def _install_idle_sigint() -> None:
try:
signal.signal(signal.SIGINT, signal.SIG_IGN)
except (OSError, ValueError):
# Some platforms (Windows in non-console mode) reject this; fine.
pass
def _install_exec_sigint() -> None:
try:
signal.signal(signal.SIGINT, signal.default_int_handler)
except (OSError, ValueError):
pass
def _start_parent_watchdog() -> None:
"""Self-terminate when the host process dies.
The main loop only exits when stdin EOFs, which only happens once user
code finishes and the next ``readline`` call returns. If the host gets
SIGKILL mid-execution (or any way that skips graceful shutdown) the
runner would otherwise outlive its parent and keep holding kernel
state. Poll ``os.getppid()`` instead and ``os._exit`` the moment we get
reparented \u2014 covers POSIX hosts. Windows has no reliable ppid
equivalent; there we still bail out on the next stdin read.
"""
if os.name != "posix":
return
original_ppid = os.getppid()
if original_ppid <= 1:
return
def watch() -> None:
while True:
try:
if os.getppid() != original_ppid:
os._exit(0)
except Exception:
return
time.sleep(10)
thread = threading.Thread(target=watch, name="omp-parent-watchdog", daemon=True)
thread.start()
# ---------------------------------------------------------------------------
# Request dispatch
# ---------------------------------------------------------------------------
def _handle_request(req: dict) -> None:
if req.get("type") == "exit":
sys.exit(0)
rid = str(req.get("id"))
code = req.get("code", "")
_STATE.current_id = rid
_STATE.cancel_requested = False
_STATE.execution_count += 1
_emit({"type": "started", "id": rid})
status: str = "ok"
cancelled = False
try:
transformed = transform_cell(code)
except SyntaxError as exc:
_emit_error(rid, exc)
_emit({
"type": "done",
"id": rid,
"status": "error",
"executionCount": _STATE.execution_count,
"cancelled": False,
})
_STATE.current_id = None
return
_install_exec_sigint()
try:
_exec_source(transformed, _STATE.user_ns)
except KeyboardInterrupt:
cancelled = True
status = "error"
_emit_error(rid, KeyboardInterrupt("Execution interrupted"))
except SystemExit:
raise
except BaseException as exc: # noqa: BLE001 - we want to surface every user error
status = "error"
_emit_error(rid, exc)
finally:
_install_idle_sigint()
try:
_flush_matplotlib_figures()
except Exception:
pass
_emit({
"type": "done",
"id": rid,
"status": status,
"executionCount": _STATE.execution_count,
"cancelled": cancelled,
})
_STATE.current_id = None
def _emit_error(rid: str, exc: BaseException) -> None:
tb_lines = traceback.format_exception(type(exc), exc, exc.__traceback__)
_emit({
"type": "error",
"id": rid,
"ename": type(exc).__name__,
"evalue": str(exc),
"traceback": [line.rstrip("\n") for line in tb_lines],
})
# ---------------------------------------------------------------------------
# Main loop
# ---------------------------------------------------------------------------
def main() -> None:
sys.stdout = _StreamProxy("stdout")
sys.stderr = _StreamProxy("stderr")
_install_idle_sigint()
_start_parent_watchdog()
stdin = sys.__stdin__
if stdin is None:
return
for raw_line in stdin:
line = raw_line.strip()
if not line:
continue
try:
req = json.loads(line)
except json.JSONDecodeError as exc:
_emit({
"type": "error",
"id": "",
"ename": "ProtocolError",
"evalue": f"Invalid JSON request: {exc}",
"traceback": [],
})
continue
try:
_handle_request(req)
except SystemExit:
return
if __name__ == "__main__":
main()
+3 -16
View File
@@ -160,19 +160,6 @@ export function resolveVenvPath(cwd: string): string | undefined {
return undefined; return undefined;
} }
/**
* Resolve the windowless Python executable (pythonw.exe) on Windows.
* Falls back to the regular Python path if pythonw.exe is not available.
*/
function resolveWindowlessPython(pythonPath: string): string {
if (process.platform !== "win32") return pythonPath;
const pythonwPath = pythonPath.replace(/python\.exe$/i, "pythonw.exe");
if (pythonwPath !== pythonPath && fs.existsSync(pythonwPath)) {
return pythonwPath;
}
return pythonPath;
}
/** /**
* Resolve Python runtime including executable path, environment, and venv detection. * Resolve Python runtime including executable path, environment, and venv detection.
*/ */
@@ -189,7 +176,7 @@ export function resolvePythonRuntime(cwd: string, baseEnv: Record<string, string
const currentPath = env[pathKey]; const currentPath = env[pathKey];
env[pathKey] = currentPath ? `${binDir}${path.delimiter}${currentPath}` : binDir; env[pathKey] = currentPath ? `${binDir}${path.delimiter}${currentPath}` : binDir;
return { return {
pythonPath: resolveWindowlessPython(pythonCandidate), pythonPath: pythonCandidate,
env, env,
venvPath, venvPath,
}; };
@@ -205,7 +192,7 @@ export function resolvePythonRuntime(cwd: string, baseEnv: Record<string, string
process.platform === "win32" ? path.join(managed.venvPath, "Scripts") : path.join(managed.venvPath, "bin"); process.platform === "win32" ? path.join(managed.venvPath, "Scripts") : path.join(managed.venvPath, "bin");
env[pathKey] = currentPath ? `${managedBin}${path.delimiter}${currentPath}` : managedBin; env[pathKey] = currentPath ? `${managedBin}${path.delimiter}${currentPath}` : managedBin;
return { return {
pythonPath: resolveWindowlessPython(managed.pythonPath), pythonPath: managed.pythonPath,
env, env,
venvPath: managed.venvPath, venvPath: managed.venvPath,
}; };
@@ -216,7 +203,7 @@ export function resolvePythonRuntime(cwd: string, baseEnv: Record<string, string
throw new Error("Python executable not found on PATH"); throw new Error("Python executable not found on PATH");
} }
return { return {
pythonPath: resolveWindowlessPython(pythonPath), pythonPath,
env, env,
}; };
} }
@@ -14,7 +14,6 @@ import { formatDuration, Snowflake, setProjectDir } from "@oh-my-pi/pi-utils";
import { $ } from "bun"; import { $ } from "bun";
import { reset as resetCapabilities } from "../../capability"; import { reset as resetCapabilities } from "../../capability";
import { clearClaudePluginRootsCache } from "../../discovery/helpers"; import { clearClaudePluginRootsCache } from "../../discovery/helpers";
import { getGatewayStatus } from "../../eval/py/gateway-coordinator";
import { loadCustomShare } from "../../export/custom-share"; import { loadCustomShare } from "../../export/custom-share";
import type { CompactOptions } from "../../extensibility/extensions/types"; import type { CompactOptions } from "../../extensibility/extensions/types";
import { import {
@@ -402,28 +401,6 @@ export class CommandController {
} }
} }
const gateway = await getGatewayStatus();
info += `\n${theme.bold("Python Gateway")}\n`;
if (gateway.active) {
info += `${theme.fg("dim", "Status:")} ${theme.fg("success", "Active (Global)")}\n`;
info += `${theme.fg("dim", "URL:")} ${gateway.url}\n`;
info += `${theme.fg("dim", "PID:")} ${gateway.pid}\n`;
if (gateway.pythonPath) {
info += `${theme.fg("dim", "Python:")} ${gateway.pythonPath}\n`;
}
if (gateway.venvPath) {
info += `${theme.fg("dim", "Venv:")} ${gateway.venvPath}\n`;
}
if (gateway.uptime !== null) {
const uptimeSec = Math.floor(gateway.uptime / 1000);
const mins = Math.floor(uptimeSec / 60);
const secs = uptimeSec % 60;
info += `${theme.fg("dim", "Uptime:")} ${mins}m ${secs}s\n`;
}
} else {
info += `${theme.fg("dim", "Status:")} ${theme.fg("dim", "Inactive")}\n`;
}
if (this.ctx.lspServers && this.ctx.lspServers.length > 0) { if (this.ctx.lspServers && this.ctx.lspServers.length > 0) {
info += `\n${theme.bold("LSP Servers")}\n`; info += `\n${theme.bold("LSP Servers")}\n`;
for (const server of this.ctx.lspServers) { for (const server of this.ctx.lspServers) {
@@ -6561,7 +6561,6 @@ export class AgentSession {
sessionId, sessionId,
kernelOwnerId: this.#evalKernelOwnerId, kernelOwnerId: this.#evalKernelOwnerId,
kernelMode: this.settings.get("python.kernelMode"), kernelMode: this.settings.get("python.kernelMode"),
useSharedGateway: this.settings.get("python.sharedGateway"),
onChunk, onChunk,
signal: abortController.signal, signal: abortController.signal,
}); });
@@ -0,0 +1,39 @@
import { describe, expect, it } from "bun:test";
import { renderKernelDisplay } from "@oh-my-pi/pi-coding-agent/eval/py/display";
describe("renderKernelDisplay (raw bundle shape)", () => {
it("renders status events without text output", async () => {
const { text, outputs } = await renderKernelDisplay({
"application/x-omp-status": { op: "find", count: 12, pattern: "foo" },
});
expect(text).toBe("");
expect(outputs).toEqual([{ type: "status", event: { op: "find", count: 12, pattern: "foo" } }]);
});
it("prefers text/markdown over text/plain", async () => {
const { text, outputs } = await renderKernelDisplay({
"text/markdown": "**bold**",
"text/plain": "bold",
});
expect(text).toBe("**bold**\n");
expect(outputs).toContainEqual({ type: "markdown" });
});
it("collects image/png alongside text/plain", async () => {
const { text, outputs } = await renderKernelDisplay({
"image/png": "base64data",
"text/plain": "<Figure>",
});
expect(text).toBe("<Figure>\n");
expect(outputs).toContainEqual({ type: "image", data: "base64data", mimeType: "image/png" });
});
it("emits json bundle and includes text/plain when present", async () => {
const { text, outputs } = await renderKernelDisplay({
"application/json": { ok: true },
"text/plain": "{ ok: true }",
});
expect(text).toBe("{ ok: true }\n");
expect(outputs).toEqual([{ type: "json", data: { ok: true } }]);
});
});
@@ -4,7 +4,6 @@ import {
disposeKernelSessionsByOwner, disposeKernelSessionsByOwner,
executePython, executePython,
} from "@oh-my-pi/pi-coding-agent/eval/py/executor"; } from "@oh-my-pi/pi-coding-agent/eval/py/executor";
import * as gatewayCoordinator from "@oh-my-pi/pi-coding-agent/eval/py/gateway-coordinator";
import type { import type {
KernelExecuteResult, KernelExecuteResult,
KernelShutdownResult, KernelShutdownResult,
@@ -320,80 +319,6 @@ describe("python executor owner cleanup", () => {
expect(kernel.shutdown).toHaveBeenCalledWith(expect.objectContaining({ timeoutMs: expect.any(Number) })); expect(kernel.shutdown).toHaveBeenCalledWith(expect.objectContaining({ timeoutMs: expect.any(Number) }));
expect(startSpy).toHaveBeenCalledTimes(1); expect(startSpy).toHaveBeenCalledTimes(1);
}); });
it("keeps local owner-cleanup disposals counted during resource-exhaustion recovery", async () => {
vi.useFakeTimers();
try {
const staleKernels = [new FakeKernel(), new FakeKernel(), new FakeKernel()];
const recoveredKernel = new FakeKernel();
const laterKernel = new FakeKernel();
const staleShutdownDeferreds = staleKernels.map(() => Promise.withResolvers<KernelShutdownResult>());
for (const [index, kernel] of staleKernels.entries()) {
kernel.shutdown = vi.fn(() => staleShutdownDeferreds[index]!.promise);
}
const shutdownSharedGatewaySpy = vi.spyOn(gatewayCoordinator, "shutdownSharedGateway").mockResolvedValue();
vi.spyOn(pythonKernel, "checkPythonKernelAvailability").mockResolvedValue({ ok: true });
const startSpy = vi.spyOn(PythonKernel, "start");
for (const kernel of staleKernels) {
startSpy.mockResolvedValueOnce(kernel as unknown as PythonKernelInstance);
}
startSpy
.mockRejectedValueOnce(new Error("EMFILE: too many open files"))
.mockResolvedValueOnce(recoveredKernel as unknown as PythonKernelInstance)
.mockResolvedValueOnce(laterKernel as unknown as PythonKernelInstance);
for (const [index] of staleKernels.entries()) {
await executePython(`print(${index})`, {
cwd: `/tmp/recovery-stale-${index}`,
sessionId: `recovery-stale-session-${index}`,
kernelMode: "session",
kernelOwnerId: "owner-a",
});
}
const ownerCleanup = disposeKernelSessionsByOwner("owner-a");
await Promise.resolve();
for (const kernel of staleKernels) {
expect(kernel.shutdown).toHaveBeenCalledWith({ timeoutMs: 2_000 });
}
vi.advanceTimersByTime(2_000);
await ownerCleanup;
await executePython("print('recovered')", {
cwd: "/tmp/recovery-after-emfile",
sessionId: "recovery-session",
kernelMode: "session",
});
expect(shutdownSharedGatewaySpy).toHaveBeenCalledTimes(1);
expect(startSpy).toHaveBeenCalledTimes(5);
expect(recoveredKernel.execute).toHaveBeenCalledTimes(1);
const blockedExecution = executePython("print('later')", {
cwd: "/tmp/recovery-after-emfile-later",
sessionId: "recovery-session-later",
kernelMode: "session",
});
await flushMicrotasks();
expect(startSpy).toHaveBeenCalledTimes(5);
expect(recoveredKernel.shutdown).not.toHaveBeenCalled();
expect(laterKernel.execute).not.toHaveBeenCalled();
staleShutdownDeferreds[0]!.resolve({ confirmed: true });
await blockedExecution;
expect(startSpy).toHaveBeenCalledTimes(6);
expect(laterKernel.execute).toHaveBeenCalledTimes(1);
for (const deferred of staleShutdownDeferreds.slice(1)) {
deferred.resolve({ confirmed: true });
}
await flushMicrotasks();
await disposeAllKernelSessions();
expect(recoveredKernel.shutdown).toHaveBeenCalledTimes(1);
expect(laterKernel.shutdown).toHaveBeenCalledTimes(1);
} finally {
vi.useRealTimers();
}
});
it("returns owner cleanup promptly but keeps retained capacity reserved until shutdown is confirmed", async () => { it("returns owner cleanup promptly but keeps retained capacity reserved until shutdown is confirmed", async () => {
vi.useFakeTimers(); vi.useFakeTimers();
try { try {
@@ -1,108 +0,0 @@
import { describe, expect, it } from "bun:test";
import {
deserializeWebSocketMessage,
type JupyterMessage,
serializeWebSocketMessage,
} from "@oh-my-pi/pi-coding-agent/eval/py/kernel";
const encoder = new TextEncoder();
function buildFrame(message: Omit<JupyterMessage, "buffers">, buffers: Uint8Array[] = []): ArrayBuffer {
const msgBytes = encoder.encode(JSON.stringify(message));
const offsetCount = 1 + buffers.length;
const headerSize = 4 + offsetCount * 4;
let totalSize = headerSize + msgBytes.length;
for (const buffer of buffers) {
totalSize += buffer.length;
}
const frame = new ArrayBuffer(totalSize);
const view = new DataView(frame);
const bytes = new Uint8Array(frame);
view.setUint32(0, offsetCount, true);
view.setUint32(4, headerSize, true);
bytes.set(msgBytes, headerSize);
let offset = headerSize + msgBytes.length;
for (let i = 0; i < buffers.length; i++) {
view.setUint32(4 + (i + 1) * 4, offset, true);
bytes.set(buffers[i], offset);
offset += buffers[i].length;
}
return frame;
}
describe("deserializeWebSocketMessage", () => {
it("parses offset tables and buffers", () => {
const message = {
channel: "iopub",
header: {
msg_id: "msg-1",
session: "session-1",
username: "omp",
date: "2024-01-01T00:00:00Z",
msg_type: "stream",
version: "5.5",
},
parent_header: {},
metadata: {},
content: { text: "hello" },
};
const buffer = new Uint8Array([1, 2, 3]);
const frame = buildFrame(message, [buffer]);
const parsed = deserializeWebSocketMessage(frame);
expect(parsed).not.toBeNull();
expect(parsed?.header.msg_id).toBe("msg-1");
expect(parsed?.content).toEqual({ text: "hello" });
expect(parsed?.buffers?.[0]).toEqual(buffer);
});
it("returns null for invalid frames", () => {
const headerSize = 8;
const bytes = encoder.encode("not-json");
const frame = new ArrayBuffer(headerSize + bytes.length);
const view = new DataView(frame);
const data = new Uint8Array(frame);
view.setUint32(0, 1, true);
view.setUint32(4, headerSize, true);
data.set(bytes, headerSize);
expect(deserializeWebSocketMessage(frame)).toBeNull();
const emptyFrame = new ArrayBuffer(4);
new DataView(emptyFrame).setUint32(0, 0, true);
expect(deserializeWebSocketMessage(emptyFrame)).toBeNull();
});
});
describe("serializeWebSocketMessage", () => {
it("round trips message payloads", () => {
const message: JupyterMessage = {
channel: "shell",
header: {
msg_id: "msg-2",
session: "session-2",
username: "omp",
date: "2024-02-01T00:00:00Z",
msg_type: "execute_request",
version: "5.5",
},
parent_header: { parent: "root" },
metadata: { tag: "meta" },
content: { code: "print('hi')" },
buffers: [new Uint8Array([9, 8, 7])],
};
const frame = serializeWebSocketMessage(message);
const parsed = deserializeWebSocketMessage(frame);
expect(parsed).not.toBeNull();
expect(parsed?.header.msg_type).toBe("execute_request");
expect(parsed?.content).toEqual({ code: "print('hi')" });
expect(parsed?.buffers?.[0]).toEqual(new Uint8Array([9, 8, 7]));
});
});
@@ -1,469 +0,0 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test";
import * as gatewayCoordinator from "@oh-my-pi/pi-coding-agent/eval/py/gateway-coordinator";
import { PythonKernel } from "@oh-my-pi/pi-coding-agent/eval/py/kernel";
import { hookFetch, TempDir } from "@oh-my-pi/pi-utils";
import type { Subprocess } from "bun";
type SpawnOptions = Bun.SpawnOptions.SpawnOptions<
Bun.SpawnOptions.Writable,
Bun.SpawnOptions.Readable,
Bun.SpawnOptions.Readable
>;
type FetchCall = { url: string; init?: RequestInit };
type FetchResponse = {
ok: boolean;
status: number;
json: () => Promise<unknown>;
text: () => Promise<string>;
};
type MockEnvironment = {
fetchCalls: FetchCall[];
spawnCalls: { cmd: string[]; options: SpawnOptions }[];
};
type MessageEventPayload = { data: ArrayBuffer };
type WebSocketHandler = (event: unknown) => void;
type WebSocketMessageHandler = (event: MessageEventPayload) => void;
class FakeWebSocket {
static OPEN = 1;
static CLOSED = 3;
static instances: FakeWebSocket[] = [];
readyState = FakeWebSocket.OPEN;
binaryType = "arraybuffer";
url: string;
sent: ArrayBuffer[] = [];
onopen: WebSocketHandler | null = null;
onerror: WebSocketHandler | null = null;
onclose: WebSocketHandler | null = null;
onmessage: WebSocketMessageHandler | null = null;
constructor(url: string) {
this.url = url;
FakeWebSocket.instances.push(this);
queueMicrotask(() => {
this.onopen?.(undefined);
});
}
send(data: ArrayBuffer): void {
this.sent.push(data);
}
close(): void {
this.readyState = FakeWebSocket.CLOSED;
this.onclose?.(undefined);
}
}
const createResponse = (options: { ok: boolean; status?: number; json?: unknown; text?: string }): FetchResponse => {
return {
ok: options.ok,
status: options.status ?? (options.ok ? 200 : 500),
json: async () => options.json ?? {},
text: async () => options.text ?? "",
};
};
const createFakeProcess = (): Subprocess => {
const exited = new Promise<number>(() => undefined);
return { pid: 999999, exited } as Subprocess;
};
const expectResolvesWithin = async <T>(promise: Promise<T>, timeoutMs: number, message: string): Promise<T> => {
let timer: ReturnType<typeof setTimeout> | undefined;
try {
return await Promise.race([
promise,
new Promise<T>((_, reject) => {
timer = setTimeout(() => reject(new Error(message)), timeoutMs);
timer.unref?.();
}),
]);
} finally {
if (timer !== undefined) clearTimeout(timer);
}
};
describe("PythonKernel gateway lifecycle", () => {
const originalWebSocket = globalThis.WebSocket;
const originalGatewayUrl = Bun.env.PI_PYTHON_GATEWAY_URL;
const originalGatewayToken = Bun.env.PI_PYTHON_GATEWAY_TOKEN;
const originalBunEnv = Bun.env.BUN_ENV;
let tempDir: TempDir;
let env: MockEnvironment;
const stubKernelRuntime = () => {
function mockSpawn(options: SpawnOptions & { cmd: string[] }): Subprocess;
function mockSpawn(cmd: string[], options?: SpawnOptions): Subprocess;
function mockSpawn(first: string[] | (SpawnOptions & { cmd: string[] }), second?: SpawnOptions): Subprocess {
if (Array.isArray(first)) {
env.spawnCalls.push({ cmd: first, options: second ?? {} });
} else {
const { cmd, ...options } = first;
env.spawnCalls.push({ cmd, options });
}
return createFakeProcess();
}
const spawnSpy = vi.spyOn(Bun, "spawn").mockImplementation(mockSpawn);
const sleepSpy = vi.spyOn(Bun, "sleep").mockImplementation(async () => undefined);
const whichSpy = vi.spyOn(Bun, "which").mockImplementation(() => "/usr/bin/python");
const executeSpy = vi.spyOn(PythonKernel.prototype, "execute").mockResolvedValue({
status: "ok",
cancelled: false,
timedOut: false,
stdinRequested: false,
});
return {
[Symbol.dispose]() {
spawnSpy.mockRestore();
sleepSpy.mockRestore();
whichSpy.mockRestore();
executeSpy.mockRestore();
},
};
};
beforeEach(() => {
tempDir = TempDir.createSync("@omp-python-kernel-");
env = { fetchCalls: [], spawnCalls: [] };
Bun.env.BUN_ENV = "test";
delete Bun.env.PI_PYTHON_GATEWAY_URL;
delete Bun.env.PI_PYTHON_GATEWAY_TOKEN;
FakeWebSocket.instances = [];
globalThis.WebSocket = FakeWebSocket as unknown as typeof WebSocket;
});
afterEach(() => {
if (tempDir) {
tempDir.removeSync();
}
if (originalBunEnv === undefined) {
delete Bun.env.BUN_ENV;
} else {
Bun.env.BUN_ENV = originalBunEnv;
}
if (originalGatewayUrl === undefined) {
delete Bun.env.PI_PYTHON_GATEWAY_URL;
} else {
Bun.env.PI_PYTHON_GATEWAY_URL = originalGatewayUrl;
}
if (originalGatewayToken === undefined) {
delete Bun.env.PI_PYTHON_GATEWAY_TOKEN;
} else {
Bun.env.PI_PYTHON_GATEWAY_TOKEN = originalGatewayToken;
}
globalThis.WebSocket = originalWebSocket;
vi.restoreAllMocks();
});
it("starts shared gateway, interrupts, and shuts down", async () => {
using _runtime = stubKernelRuntime();
vi.spyOn(gatewayCoordinator, "acquireSharedGateway").mockResolvedValue({
url: "http://127.0.0.1:9999",
isShared: true,
});
using _hook = hookFetch((input, init) => {
const url = String(input);
env.fetchCalls.push({ url, init });
if (url.endsWith("/api/kernels") && init?.method === "POST") {
return createResponse({ ok: true, json: { id: "kernel-123" } }) as unknown as Response;
}
return createResponse({ ok: true }) as unknown as Response;
});
const kernel = await PythonKernel.start({ cwd: tempDir.path() });
expect(env.fetchCalls.some(call => call.url.endsWith("/api/kernels") && call.init?.method === "POST")).toBe(true);
await kernel.interrupt();
expect(env.fetchCalls.some(call => call.url.includes("/interrupt") && call.init?.method === "POST")).toBe(true);
await kernel.shutdown();
expect(env.fetchCalls.some(call => call.init?.method === "DELETE")).toBe(true);
expect(kernel.isAlive()).toBe(false);
});
it("aborts stalled startup after websocket connect and cleans up the kernel", async () => {
vi.spyOn(gatewayCoordinator, "acquireSharedGateway").mockResolvedValue({
url: "http://127.0.0.1:9999",
isShared: true,
});
let executeCallCount = 0;
const preludeStarted = Promise.withResolvers<void>();
vi.spyOn(PythonKernel.prototype, "execute").mockImplementation(async (_code, options) => {
executeCallCount += 1;
if (executeCallCount === 1) {
return { status: "ok", cancelled: false, timedOut: false, stdinRequested: false };
}
preludeStarted.resolve();
return await new Promise((_, reject) => {
const onAbort = () => {
const reason = options?.signal?.reason;
reject(reason instanceof Error ? reason : new Error("Python kernel startup aborted"));
};
if (options?.signal?.aborted) {
onAbort();
return;
}
options?.signal?.addEventListener("abort", onAbort, { once: true });
});
});
using _hook = hookFetch((input, init) => {
const url = String(input);
env.fetchCalls.push({ url, init });
if (url.endsWith("/api/kernels") && init?.method === "POST") {
return createResponse({ ok: true, json: { id: "kernel-stalled" } }) as unknown as Response;
}
return createResponse({ ok: true }) as unknown as Response;
});
const abortController = new AbortController();
const startPromise = PythonKernel.start({ cwd: tempDir.path(), signal: abortController.signal });
await preludeStarted.promise;
abortController.abort(new Error("cancel startup"));
const pending = Symbol("pending");
const settled = await Promise.race([
startPromise.then(
() => "resolved",
error => error,
),
Bun.sleep(50).then(() => pending),
]);
expect(settled).toBeInstanceOf(Error);
expect(settled).not.toBe(pending);
expect((settled as Error).message).toContain("cancel startup");
await expect(startPromise).rejects.toThrow("cancel startup");
expect(
env.fetchCalls.some(
call => call.url.endsWith("/api/kernels/kernel-stalled") && call.init?.method === "DELETE",
),
).toBe(true);
expect(FakeWebSocket.instances.at(-1)?.readyState).toBe(FakeWebSocket.CLOSED);
});
it("preserves timeout classification when startup environment initialization is cancelled", async () => {
vi.spyOn(gatewayCoordinator, "acquireSharedGateway").mockResolvedValue({
url: "http://127.0.0.1:9999",
isShared: true,
});
vi.spyOn(PythonKernel.prototype, "execute").mockResolvedValue({
status: "ok",
cancelled: true,
timedOut: true,
stdinRequested: false,
});
using _hook = hookFetch((input, init) => {
const url = String(input);
env.fetchCalls.push({ url, init });
if (url.endsWith("/api/kernels") && init?.method === "POST") {
return createResponse({ ok: true, json: { id: "kernel-init-timeout" } }) as unknown as Response;
}
return createResponse({ ok: true }) as unknown as Response;
});
await expect(PythonKernel.start({ cwd: tempDir.path() })).rejects.toMatchObject({
name: "TimeoutError",
message: "Failed to initialize Python kernel environment",
});
expect(
env.fetchCalls.some(
call => call.url.endsWith("/api/kernels/kernel-init-timeout") && call.init?.method === "DELETE",
),
).toBe(true);
expect(FakeWebSocket.instances.at(-1)?.readyState).toBe(FakeWebSocket.CLOSED);
});
it("preserves timeout classification when startup prelude execution is cancelled", async () => {
vi.spyOn(gatewayCoordinator, "acquireSharedGateway").mockResolvedValue({
url: "http://127.0.0.1:9999",
isShared: true,
});
let executeCallCount = 0;
vi.spyOn(PythonKernel.prototype, "execute").mockImplementation(async () => {
executeCallCount += 1;
if (executeCallCount === 1) {
return { status: "ok", cancelled: false, timedOut: false, stdinRequested: false };
}
return { status: "ok", cancelled: true, timedOut: true, stdinRequested: false };
});
using _hook = hookFetch((input, init) => {
const url = String(input);
env.fetchCalls.push({ url, init });
if (url.endsWith("/api/kernels") && init?.method === "POST") {
return createResponse({ ok: true, json: { id: "kernel-prelude-timeout" } }) as unknown as Response;
}
return createResponse({ ok: true }) as unknown as Response;
});
await expect(PythonKernel.start({ cwd: tempDir.path() })).rejects.toMatchObject({
name: "TimeoutError",
message: "Failed to initialize Python kernel prelude",
});
expect(
env.fetchCalls.some(
call => call.url.endsWith("/api/kernels/kernel-prelude-timeout") && call.init?.method === "DELETE",
),
).toBe(true);
expect(FakeWebSocket.instances.at(-1)?.readyState).toBe(FakeWebSocket.CLOSED);
});
it("throws when shared gateway kernel creation never succeeds", async () => {
using _runtime = stubKernelRuntime();
vi.spyOn(gatewayCoordinator, "acquireSharedGateway").mockResolvedValue({
url: "http://127.0.0.1:9999",
isShared: true,
});
using _hook = hookFetch((input, init) => {
const url = String(input);
env.fetchCalls.push({ url, init });
if (url.endsWith("/api/kernels") && init?.method === "POST") {
return createResponse({ ok: false, status: 503, text: "oops" }) as unknown as Response;
}
return createResponse({ ok: true }) as unknown as Response;
});
await expect(PythonKernel.start({ cwd: tempDir.path() })).rejects.toThrow(
"Failed to create kernel on shared gateway",
);
});
it("treats initial 404 and 410 shutdown responses as confirmed", async () => {
using _runtime = stubKernelRuntime();
vi.spyOn(gatewayCoordinator, "acquireSharedGateway").mockResolvedValue({
url: "http://127.0.0.1:9999",
isShared: true,
});
for (const status of [404, 410]) {
let deleteCalls = 0;
using _hook = hookFetch((input, init) => {
const url = String(input);
env.fetchCalls.push({ url, init });
if (url.endsWith("/api/kernels") && init?.method === "POST") {
return createResponse({ ok: true, json: { id: `kernel-missing-${status}` } }) as unknown as Response;
}
if (url.endsWith(`/api/kernels/kernel-missing-${status}`) && init?.method === "DELETE") {
deleteCalls += 1;
return createResponse({ ok: false, status, text: "gone" }) as unknown as Response;
}
return createResponse({ ok: true }) as unknown as Response;
});
const kernel = await PythonKernel.start({ cwd: tempDir.path() });
await expect(kernel.shutdown()).resolves.toEqual({ confirmed: true });
expect(deleteCalls).toBe(1);
expect(kernel.isAlive()).toBe(false);
expect(FakeWebSocket.instances.at(-1)?.readyState).toBe(FakeWebSocket.CLOSED);
await expect(kernel.shutdown()).resolves.toEqual({ confirmed: true });
expect(deleteCalls).toBe(1);
}
});
it("returns unconfirmed when shutdown times out and can confirm on retry", async () => {
using _runtime = stubKernelRuntime();
vi.spyOn(gatewayCoordinator, "acquireSharedGateway").mockResolvedValue({
url: "http://127.0.0.1:9999",
isShared: true,
});
let deleteCalls = 0;
const firstDeleteStarted = Promise.withResolvers<void>();
const firstDeleteAborted = Promise.withResolvers<void>();
using _hook = hookFetch((input, init) => {
const url = String(input);
env.fetchCalls.push({ url, init });
if (url.endsWith("/api/kernels") && init?.method === "POST") {
return createResponse({ ok: true, json: { id: "kernel-shutdown-timeout" } }) as unknown as Response;
}
if (url.endsWith("/api/kernels/kernel-shutdown-timeout") && init?.method === "DELETE") {
deleteCalls += 1;
if (deleteCalls === 1) {
firstDeleteStarted.resolve();
return new Promise<Response>((_, reject) => {
const abortSignal = init.signal;
if (!abortSignal) return;
const rejectOnAbort = () => {
firstDeleteAborted.resolve();
const reason = abortSignal.reason;
reject(reason instanceof Error ? reason : new Error("Python kernel shutdown timed out"));
};
if (abortSignal.aborted) {
rejectOnAbort();
return;
}
abortSignal.addEventListener("abort", rejectOnAbort, { once: true });
});
}
return createResponse({ ok: false, status: 404, text: "gone" }) as unknown as Response;
}
return createResponse({ ok: true }) as unknown as Response;
});
const kernel = await PythonKernel.start({ cwd: tempDir.path() });
const shutdownPromise = kernel.shutdown({ timeoutMs: 25 });
await expectResolvesWithin(firstDeleteStarted.promise, 250, "kernel shutdown never issued a delete request");
await expectResolvesWithin(
firstDeleteAborted.promise,
500,
"timed out waiting for the first delete request to abort",
);
await expect(
expectResolvesWithin(shutdownPromise, 500, "kernel shutdown did not settle after timing out"),
).resolves.toEqual({
confirmed: false,
});
expect(kernel.isAlive()).toBe(false);
expect(FakeWebSocket.instances.at(-1)?.readyState).toBe(FakeWebSocket.CLOSED);
await expect(kernel.shutdown()).resolves.toEqual({ confirmed: true });
expect(deleteCalls).toBe(2);
});
it("does not throw when shutdown API fails", async () => {
using _runtime = stubKernelRuntime();
vi.spyOn(gatewayCoordinator, "acquireSharedGateway").mockResolvedValue({
url: "http://127.0.0.1:9999",
isShared: true,
});
using _hook = hookFetch((input, init) => {
const url = String(input);
env.fetchCalls.push({ url, init });
if (url.endsWith("/api/kernels") && init?.method === "POST") {
return createResponse({ ok: true, json: { id: "kernel-456" } }) as unknown as Response;
}
if (init?.method === "DELETE") {
throw new Error("delete failed");
}
return createResponse({ ok: true }) as unknown as Response;
});
const kernel = await PythonKernel.start({ cwd: tempDir.path() });
await expect(kernel.shutdown()).resolves.toEqual({ confirmed: false });
});
});
@@ -1,377 +0,0 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test";
import { type KernelDisplayOutput, PythonKernel } from "@oh-my-pi/pi-coding-agent/eval/py/kernel";
import { PYTHON_PRELUDE } from "@oh-my-pi/pi-coding-agent/eval/py/prelude";
import { hookFetch } from "@oh-my-pi/pi-utils";
type JupyterMessage = {
channel: string;
header: {
msg_id: string;
session: string;
username: string;
date: string;
msg_type: string;
version: string;
};
parent_header: Record<string, unknown>;
metadata: Record<string, unknown>;
content: Record<string, unknown>;
buffers?: Uint8Array[];
};
const textEncoder = new TextEncoder();
const textDecoder = new TextDecoder();
function encodeMessage(msg: JupyterMessage): ArrayBuffer {
const msgText = JSON.stringify({
channel: msg.channel,
header: msg.header,
parent_header: msg.parent_header,
metadata: msg.metadata,
content: msg.content,
});
const msgBytes = textEncoder.encode(msgText);
const buffers = msg.buffers ?? [];
const offsetCount = 1 + buffers.length;
const headerSize = 4 + offsetCount * 4;
let totalSize = headerSize + msgBytes.length;
for (const buffer of buffers) {
totalSize += buffer.length;
}
const result = new ArrayBuffer(totalSize);
const view = new DataView(result);
const bytes = new Uint8Array(result);
view.setUint32(0, offsetCount, true);
let offset = headerSize;
view.setUint32(4, offset, true);
bytes.set(msgBytes, offset);
offset += msgBytes.length;
buffers.forEach((buffer, index) => {
view.setUint32(4 + (index + 1) * 4, offset, true);
bytes.set(buffer, offset);
offset += buffer.length;
});
return result;
}
function decodeMessage(data: ArrayBuffer): JupyterMessage {
const view = new DataView(data);
const offsetCount = view.getUint32(0, true);
const offsets: number[] = [];
for (let i = 0; i < offsetCount; i++) {
offsets.push(view.getUint32(4 + i * 4, true));
}
const msgStart = offsets[0];
const msgEnd = offsets.length > 1 ? offsets[1] : data.byteLength;
const msgBytes = new Uint8Array(data, msgStart, msgEnd - msgStart);
const msgText = textDecoder.decode(msgBytes);
return JSON.parse(msgText) as JupyterMessage;
}
function sendOkExecution(ws: FakeWebSocket, msgId: string, executionCount = 1) {
const reply: JupyterMessage = {
channel: "shell",
header: {
msg_id: `reply-${msgId}`,
session: "session",
username: "omp",
date: new Date().toISOString(),
msg_type: "execute_reply",
version: "5.5",
},
parent_header: { msg_id: msgId },
metadata: {},
content: { status: "ok", execution_count: executionCount },
};
const status: JupyterMessage = {
channel: "iopub",
header: {
msg_id: `status-${msgId}`,
session: "session",
username: "omp",
date: new Date().toISOString(),
msg_type: "status",
version: "5.5",
},
parent_header: { msg_id: msgId },
metadata: {},
content: { execution_state: "idle" },
};
ws.onmessage?.({ data: encodeMessage(reply) });
ws.onmessage?.({ data: encodeMessage(status) });
}
class FakeWebSocket {
static OPEN = 1;
static CLOSED = 3;
static lastInstance: FakeWebSocket | null = null;
readyState = FakeWebSocket.OPEN;
binaryType = "arraybuffer";
onopen?: () => void;
onmessage?: (event: { data: ArrayBuffer }) => void;
onerror?: (event: unknown) => void;
onclose?: () => void;
readonly url: string;
readonly sent: (ArrayBuffer | string)[] = [];
private handleSend: ((data: ArrayBuffer | string) => void) | null = null;
private pendingMessages: (ArrayBuffer | string)[] = [];
constructor(url: string) {
this.url = url;
FakeWebSocket.lastInstance = this;
queueMicrotask(() => this.onopen?.());
}
setSendHandler(handler: (data: ArrayBuffer | string) => void) {
this.handleSend = handler;
for (const msg of this.pendingMessages) {
handler(msg);
}
this.pendingMessages = [];
}
send(data: ArrayBuffer | string) {
this.sent.push(data);
if (this.handleSend) {
this.handleSend(data);
} else {
this.pendingMessages.push(data);
}
}
close() {
this.readyState = FakeWebSocket.CLOSED;
this.onclose?.();
}
}
describe("PythonKernel (external gateway)", () => {
const originalEnv = { ...Bun.env };
const originalWebSocket = globalThis.WebSocket;
beforeEach(() => {
Bun.env.PI_PYTHON_GATEWAY_URL = "http://gateway.test";
Bun.env.PI_PYTHON_SKIP_CHECK = "1";
globalThis.WebSocket = FakeWebSocket as unknown as typeof WebSocket;
});
afterEach(() => {
for (const key of Object.keys(Bun.env)) {
if (!(key in originalEnv)) {
delete Bun.env[key];
}
}
for (const [key, value] of Object.entries(originalEnv)) {
Bun.env[key] = value;
}
globalThis.WebSocket = originalWebSocket;
FakeWebSocket.lastInstance = null;
vi.restoreAllMocks();
});
it("executes code via websocket stream and display data", async () => {
const fetchMock = vi.fn(async (url: string, init?: RequestInit) => {
if (url.endsWith("/api/kernels") && init?.method === "POST") {
return new Response(JSON.stringify({ id: "kernel-1" }), { status: 201 });
}
if (url.includes("/api/kernels/") && init?.method === "DELETE") {
return new Response("", { status: 204 });
}
return new Response("", { status: 200 });
});
using _hook = hookFetch((input, init) => fetchMock(String(input), init));
let initSeen = false;
let preludeSeen = false;
const kernelPromise = PythonKernel.start({ cwd: "/" });
await Bun.sleep(10);
const ws = FakeWebSocket.lastInstance;
if (!ws) throw new Error("WebSocket not initialized");
ws.setSendHandler(data => {
const msg = typeof data === "string" ? (JSON.parse(data) as JupyterMessage) : decodeMessage(data);
const code = String(msg.content.code ?? "");
if (!initSeen) {
// First execution is kernel environment init
initSeen = true;
sendOkExecution(ws, msg.header.msg_id);
return;
}
if (!preludeSeen) {
expect(code).toBe(PYTHON_PRELUDE);
preludeSeen = true;
sendOkExecution(ws, msg.header.msg_id);
return;
}
if (code === "print('hello')") {
const stream: JupyterMessage = {
channel: "iopub",
header: {
msg_id: "stream-1",
session: "session",
username: "omp",
date: new Date().toISOString(),
msg_type: "stream",
version: "5.5",
},
parent_header: { msg_id: msg.header.msg_id },
metadata: {},
content: { text: "hello\n" },
};
const display: JupyterMessage = {
channel: "iopub",
header: {
msg_id: "display-1",
session: "session",
username: "omp",
date: new Date().toISOString(),
msg_type: "execute_result",
version: "5.5",
},
parent_header: { msg_id: msg.header.msg_id },
metadata: {},
content: {
data: {
"text/plain": "result",
"application/json": { answer: 42 },
},
},
};
const reply: JupyterMessage = {
channel: "shell",
header: {
msg_id: "reply-2",
session: "session",
username: "omp",
date: new Date().toISOString(),
msg_type: "execute_reply",
version: "5.5",
},
parent_header: { msg_id: msg.header.msg_id },
metadata: {},
content: { status: "ok", execution_count: 2 },
};
const status: JupyterMessage = {
channel: "iopub",
header: {
msg_id: "status-2",
session: "session",
username: "omp",
date: new Date().toISOString(),
msg_type: "status",
version: "5.5",
},
parent_header: { msg_id: msg.header.msg_id },
metadata: {},
content: { execution_state: "idle" },
};
ws.onmessage?.({ data: encodeMessage(stream) });
ws.onmessage?.({ data: encodeMessage(display) });
ws.onmessage?.({ data: encodeMessage(reply) });
ws.onmessage?.({ data: encodeMessage(status) });
return;
}
sendOkExecution(ws, msg.header.msg_id);
});
const kernel = await kernelPromise;
const chunks: string[] = [];
const displays: KernelDisplayOutput[] = [];
const result = await kernel.execute("print('hello')", {
onChunk: text => {
chunks.push(text);
},
onDisplay: output => {
displays.push(output);
},
});
expect(result.status).toBe("ok");
expect(chunks.join("")).toContain("hello");
expect(chunks.join("")).toContain("result");
expect(displays).toEqual([{ type: "json", data: { answer: 42 } }]);
const shutdown = await kernel.shutdown();
expect(shutdown).toEqual({ confirmed: true });
expect(fetchMock).toHaveBeenCalledWith("http://gateway.test/api/kernels/kernel-1", {
method: "DELETE",
headers: {},
});
});
it("returns an unconfirmed shutdown result when kernel deletion is not acknowledged", async () => {
let deleteAttempts = 0;
const fetchMock = vi.fn(async (url: string, init?: RequestInit) => {
if (url.endsWith("/api/kernels") && init?.method === "POST") {
return new Response(JSON.stringify({ id: "kernel-delete-failure" }), { status: 201 });
}
if (url.includes("/api/kernels/") && init?.method === "DELETE") {
deleteAttempts += 1;
if (deleteAttempts === 1) {
return new Response("delete failed", { status: 500, statusText: "Server Error" });
}
return new Response("already gone", { status: 404, statusText: "Not Found" });
}
return new Response("", { status: 200 });
});
using _hook = hookFetch((input, init) => fetchMock(String(input), init));
const kernelPromise = PythonKernel.start({ cwd: "/" });
await Bun.sleep(10);
const ws = FakeWebSocket.lastInstance;
if (!ws) throw new Error("WebSocket not initialized");
ws.setSendHandler(data => {
const msg = typeof data === "string" ? (JSON.parse(data) as JupyterMessage) : decodeMessage(data);
sendOkExecution(ws, msg.header.msg_id);
});
const kernel = await kernelPromise;
await expect(kernel.shutdown()).resolves.toEqual({ confirmed: false });
expect(deleteAttempts).toBe(1);
await expect(kernel.shutdown()).resolves.toEqual({ confirmed: true });
expect(deleteAttempts).toBe(2);
await expect(kernel.shutdown()).resolves.toEqual({ confirmed: true });
expect(deleteAttempts).toBe(2);
});
it("initializes the IPython prelude", async () => {
const fetchMock = vi.fn(async (url: string, init?: RequestInit) => {
if (url.endsWith("/api/kernels") && init?.method === "POST") {
return new Response(JSON.stringify({ id: "kernel-3" }), { status: 201 });
}
return new Response("", { status: 200 });
});
using _hook = hookFetch((input, init) => fetchMock(String(input), init));
let initSeen = false;
let preludeSeen = false;
const kernelPromise = PythonKernel.start({ cwd: "/" });
await Bun.sleep(10);
const ws = FakeWebSocket.lastInstance;
if (!ws) throw new Error("WebSocket not initialized");
ws.setSendHandler(data => {
const msg = typeof data === "string" ? (JSON.parse(data) as JupyterMessage) : decodeMessage(data);
const code = String(msg.content.code ?? "");
if (!initSeen) {
initSeen = true;
sendOkExecution(ws, msg.header.msg_id);
return;
}
if (!preludeSeen) {
expect(code).toBe(PYTHON_PRELUDE);
preludeSeen = true;
}
sendOkExecution(ws, msg.header.msg_id);
});
const kernel = await kernelPromise;
expect(kernel.isAlive()).toBe(true);
await kernel.shutdown();
});
});
// TODO: add coverage for gateway process exit handling once PythonKernel exposes a test hook.
@@ -0,0 +1,106 @@
/**
* End-to-end exercise of the new subprocess-backed Python runner.
*
* Gated by `PI_PYTHON_INTEGRATION=1` so CI without a real Python interpreter
* (or sandboxes where subprocess spawning is restricted) does not fail.
*/
import { afterEach, describe, expect, it } from "bun:test";
import * as path from "node:path";
import { disposeAllKernelSessions, executePythonWithKernel } from "@oh-my-pi/pi-coding-agent/eval/py/executor";
import { PythonKernel } from "@oh-my-pi/pi-coding-agent/eval/py/kernel";
import { TempDir } from "@oh-my-pi/pi-utils";
const SHOULD_RUN = Bun.env.PI_PYTHON_INTEGRATION === "1";
describe.skipIf(!SHOULD_RUN)("python runner subprocess", () => {
afterEach(async () => {
await disposeAllKernelSessions();
});
it("streams stdout chunks as they are produced", async () => {
using tempDir = TempDir.createSync("@python-runner-stream-");
const kernel = await PythonKernel.start({ cwd: tempDir.path() });
try {
const chunks: string[] = [];
const result = await executePythonWithKernel(
kernel,
["import sys", "for i in range(5):", " print(i, flush=True)"].join("\n"),
{
onChunk: chunk => {
chunks.push(chunk);
},
},
);
expect(result.exitCode).toBe(0);
// 5 lines * (digit + newline) → at least 5 distinct chunks once printed.
const text = chunks.join("");
expect(text).toContain("0\n");
expect(text).toContain("4\n");
} finally {
await kernel.shutdown();
}
});
it("cancels a long sleep via SIGINT within 500ms", async () => {
using tempDir = TempDir.createSync("@python-runner-cancel-");
const kernel = await PythonKernel.start({ cwd: tempDir.path() });
try {
const start = Date.now();
const ac = new AbortController();
const pending = executePythonWithKernel(kernel, "import time\ntime.sleep(30)", {
signal: ac.signal,
});
setTimeout(() => ac.abort(new DOMException("user cancelled", "AbortError")), 50);
const result = await pending;
const elapsed = Date.now() - start;
expect(result.cancelled).toBe(true);
expect(elapsed).toBeLessThan(2_000);
// Kernel must survive cancellation and remain usable.
const next = await executePythonWithKernel(kernel, "print('alive')");
expect(next.exitCode).toBe(0);
expect(next.output).toContain("alive");
} finally {
await kernel.shutdown();
}
});
it("preserves user namespace across calls", async () => {
using tempDir = TempDir.createSync("@python-runner-session-");
const kernel = await PythonKernel.start({ cwd: tempDir.path() });
try {
await executePythonWithKernel(kernel, "x = 41");
const result = await executePythonWithKernel(kernel, "x + 1");
expect(result.exitCode).toBe(0);
expect(result.output).toContain("42");
} finally {
await kernel.shutdown();
}
});
it("emits an error frame when user code raises", async () => {
using tempDir = TempDir.createSync("@python-runner-error-");
const kernel = await PythonKernel.start({ cwd: tempDir.path() });
try {
const result = await executePythonWithKernel(kernel, "raise ValueError('boom')");
expect(result.exitCode).toBe(1);
expect(result.output).toContain("ValueError");
expect(result.output).toContain("boom");
} finally {
await kernel.shutdown();
}
});
it("translates %pwd magic to the user namespace", async () => {
using tempDir = TempDir.createSync("@python-runner-magic-");
const kernel = await PythonKernel.start({ cwd: tempDir.path() });
try {
const result = await executePythonWithKernel(kernel, "%pwd");
expect(result.exitCode).toBe(0);
// %pwd returns the cwd string, which becomes the last-expression result.
// On macOS, the OS may resolve /var to /private/var, so check by basename.
expect(result.output).toContain(path.basename(tempDir.path()));
} finally {
await kernel.shutdown();
}
});
});