feat(coding-agent): added support for Julia and display language icons in code cells

- Added Julia language support to the theme symbol maps.
- Enabled language icons in code cell headers for the eval tool renderer.
This commit is contained in:
can1357
2026-06-21 19:42:01 +02:00
parent 1f3f3cf5d1
commit 33e2594f03
21 changed files with 1115 additions and 177 deletions
+1 -1
View File
@@ -7,7 +7,7 @@
- Added a Julia eval backend (`language: "jl"`): a persistent subprocess Julia kernel that shares the standard agent tool bridge and evaluation model, supporting Julia REPL-style auto-display of results in a persistent session
- Added configuration options `eval.rb` and `eval.jl` to enable or disable the Ruby and Julia backends, as well as `ruby.interpreter` and `julia.interpreter` paths for explicit runtime control
- Added a Ruby eval backend (`language: "rb"`): a persistent subprocess Ruby kernel modeled on the Python kernel, speaking the same NDJSON protocol with isolated frame/stdout channels, clean SIGINT cancellation that preserves kernel state, and per-owner kernel cleanup. Ships the full prelude helper surface (`display`/`read`/`write`/`append`/`tree`/`diff`/`env`/`output`/`sort`/`uniq`/`counter`, the `tool.<name>` bridge proxy, and `completion`/`agent`/`parallel`/`pipeline`/`log`/`phase`/`budget`) over the shared loopback tool bridge, honors `local://` roots, and auto-displays the last expression IRB-style (suppressing assignments/definitions). Gated by the `eval.rb` setting and `PI_RB` env flag, with an optional `ruby.interpreter` override.
- Added a Julia eval backend (`language: "jl"`): a persistent subprocess Julia kernel with the same NDJSON/result-display/tool-bridge model as the Python and Ruby eval runtimes, including `eval.jl` / `PI_JL` gating, `julia.interpreter` override support, persistent per-session state, clean cancellation, owner-based kernel cleanup, `local://` path roots, and the shared prelude helper surface for file I/O, tool calls, subagents, concurrency helpers, and budget access.
- Added a Julia eval backend (`language: "jl"`): a persistent subprocess Julia kernel with the same NDJSON/result-display/tool-bridge model as the Python and Ruby eval runtimes, including `eval.jl` / `PI_JL` gating, `julia.interpreter` override support, persistent per-session state, clean cancellation, owner-based kernel cleanup, `local://` path roots, and the current Julia prelude helper surface (`display`/`read`/`write`/`append`/`tree`/`diff`/`env`/`output`, the `tool.<name>` bridge proxy, and bridge-backed `completion`/`agent`/`parallel`/`pipeline`/`log`/`phase`/`budget`, including structured `schema` parsing and `return_handle` DAG shaping).
### Changed
@@ -122,7 +122,7 @@ export const TAB_GROUPS: Record<SettingTab, readonly string[]> = {
context: ["General", "Compaction", "Rules (TTSR)", "Experimental"],
memory: ["General", "Auto-Learn", "Mnemopi", "Hindsight"],
files: ["Editing", "Reading", "Read Summaries", "LSP"],
shell: ["Bash", "Eval & Python"],
shell: ["Bash", "Eval & Runtimes"],
tools: [
"Available Tools",
"Todos",
@@ -2969,7 +2969,7 @@ export const SETTINGS_SCHEMA = {
default: true,
ui: {
tab: "shell",
group: "Eval & Python",
group: "Eval & Runtimes",
label: "Python Eval Backend",
description: "Allow the eval tool to dispatch Python cells to the IPython kernel",
},
@@ -2980,7 +2980,7 @@ export const SETTINGS_SCHEMA = {
default: true,
ui: {
tab: "shell",
group: "Eval & Python",
group: "Eval & Runtimes",
label: "JavaScript Eval Backend",
description: "Allow the eval tool to dispatch JavaScript cells to the in-process runtime",
},
@@ -2991,7 +2991,7 @@ export const SETTINGS_SCHEMA = {
default: true,
ui: {
tab: "shell",
group: "Eval & Python",
group: "Eval & Runtimes",
label: "Ruby Eval Backend",
description: "Allow the eval tool to dispatch Ruby cells to the persistent Ruby kernel",
},
@@ -3002,20 +3002,20 @@ export const SETTINGS_SCHEMA = {
default: true,
ui: {
tab: "shell",
group: "Eval & Python",
group: "Eval & Runtimes",
label: "Julia Eval Backend",
description: "Allow the eval tool to dispatch Julia cells to the persistent Julia kernel",
},
},
// Python kernel knobs (consumed by the eval py backend and the /python slash command)
// Runtime knobs (consumed by eval backends and the /python slash command)
"python.kernelMode": {
type: "enum",
values: ["session", "per-call"] as const,
default: "session",
ui: {
tab: "shell",
group: "Eval & Python",
group: "Eval & Runtimes",
label: "Python Kernel Mode",
description: "Keep the IPython kernel alive across eval calls or start fresh each time",
},
@@ -3025,7 +3025,7 @@ export const SETTINGS_SCHEMA = {
default: "",
ui: {
tab: "shell",
group: "Eval & Python",
group: "Eval & Runtimes",
label: "Python Interpreter",
description:
"Optional path to an exact Python executable. When set, automatic Python runtime discovery is skipped.",
@@ -3036,7 +3036,7 @@ export const SETTINGS_SCHEMA = {
default: "",
ui: {
tab: "shell",
group: "Eval & Python",
group: "Eval & Runtimes",
label: "Ruby Interpreter",
description:
"Optional path to an exact Ruby executable. When set, automatic Ruby runtime discovery is skipped.",
@@ -3047,7 +3047,7 @@ export const SETTINGS_SCHEMA = {
default: "",
ui: {
tab: "shell",
group: "Eval & Python",
group: "Eval & Runtimes",
label: "Julia Interpreter",
description:
"Optional path to an exact Julia executable. When set, automatic Julia runtime discovery is skipped.",
@@ -0,0 +1,77 @@
import { afterEach, describe, expect, it } from "bun:test";
import * as path from "node:path";
import { $which, TempDir } from "@oh-my-pi/pi-utils";
import { disposeJuliaKernelSessionsByOwner, executeJulia } from "../jl/executor";
const HAS_JULIA = Boolean($which("julia"));
const OWNER_ID = "julia-prelude-tests";
describe.skipIf(!HAS_JULIA)("eval Julia prelude helpers", () => {
afterEach(async () => {
await disposeJuliaKernelSessionsByOwner(OWNER_ID);
});
it("supports tree keyword options and unified diff", async () => {
using tempDir = TempDir.createSync("@omp-eval-julia-helpers-");
await Bun.write(path.join(tempDir.path(), "a.txt"), "same\nold\n");
await Bun.write(path.join(tempDir.path(), "b.txt"), "same\nnew\n");
await Bun.write(path.join(tempDir.path(), "dir", "child.txt"), "child");
const result = await executeJulia(
`
d = diff("a.txt", "b.txt")
println("DIFF_DELETE=", occursin("-old", d))
println("DIFF_ADD=", occursin("+new", d))
t = tree(".", max_depth=2)
println("TREE_CHILD=", occursin("child.txt", t))
nothing
`,
{
cwd: tempDir.path(),
sessionId: `julia-prelude-diff:${crypto.randomUUID()}`,
kernelOwnerId: OWNER_ID,
reset: true,
},
);
expect(result.exitCode).toBe(0);
expect(result.output).toContain("DIFF_DELETE=true");
expect(result.output).toContain("DIFF_ADD=true");
expect(result.output).toContain("TREE_CHILD=true");
});
it("supports output ranges, JSON queries, metadata, and ANSI stripping", async () => {
using tempDir = TempDir.createSync("@omp-eval-julia-output-");
const artifactsDir = path.join(tempDir.path(), "session-artifacts");
await Bun.write(path.join(artifactsDir, "alpha.md"), "one\ntwo\nthree\nfour");
await Bun.write(path.join(artifactsDir, "json.md"), JSON.stringify({ items: [{ name: "a" }, { name: "b" }] }));
await Bun.write(path.join(artifactsDir, "ansi.md"), "\u001b[31mred\u001b[0m");
const result = await executeJulia(
`
println("RANGE=", replace(output("alpha", offset=2, limit=2), "\\n" => "|"))
println("QUERY=", output("json", query=".items[1].name"))
println("STRIPPED=", output("ansi", format="stripped"))
meta = output("alpha", format="json")
println("META=", meta["id"], ":", meta["char_count"] > 0)
multi = output("alpha", "json")
println("MULTI=", length(multi), ":", multi[1]["id"], ":", multi[2]["id"])
nothing
`,
{
cwd: tempDir.path(),
artifactsDir,
sessionId: `julia-prelude-output:${crypto.randomUUID()}`,
kernelOwnerId: OWNER_ID,
reset: true,
},
);
expect(result.exitCode).toBe(0);
expect(result.output).toContain("RANGE=two|three");
expect(result.output).toContain('QUERY="b"');
expect(result.output).toContain("STRIPPED=red");
expect(result.output).toContain("META=alpha:true");
expect(result.output).toContain("MULTI=2:alpha:json");
});
});
+137 -43
View File
@@ -56,11 +56,20 @@ export interface JuliaResult {
stdinRequested: boolean;
}
interface JuliaSession {
interface JuliaSessionOwners {
ownerIds: Set<string>;
hasFallbackOwner: boolean;
}
interface JuliaSession extends JuliaSessionOwners {
sessionKey: string;
sessionId: string;
cwd: string;
kernel: JuliaKernel;
owners: Set<string>;
}
interface StartingJuliaSession extends JuliaSessionOwners {
promise: Promise<JuliaSession>;
}
class JuliaExecutionCancelledError extends Error {
@@ -71,7 +80,7 @@ class JuliaExecutionCancelledError extends Error {
}
const sessions = new Map<string, JuliaSession>();
const startingSessions = new Map<string, Promise<JuliaSession>>();
const startingSessions = new Map<string, StartingJuliaSession>();
const resettingSessions = new Map<string, Promise<void>>();
function normalizeSessionCwd(cwd: string): string {
@@ -239,6 +248,7 @@ function buildKernelEnv(options: {
}
async function startKernel(cwd: string, options: JuliaExecutorOptions): Promise<JuliaKernel> {
requireRemainingTimeoutMs(options.deadlineMs);
const env: Record<string, string | undefined> = {};
const patch = buildKernelEnv(options);
if (patch) {
@@ -256,11 +266,18 @@ async function startKernel(cwd: string, options: JuliaExecutorOptions): Promise<
});
}
function attachOwner(session: JuliaSession, sessionId: string, ownerId: string | undefined): void {
if (ownerId) {
session.owners.add(ownerId);
} else {
session.owners.add(`unmanaged:${sessionId}`);
function attachOwner(session: JuliaSessionOwners, sessionId: string, ownerId: string | undefined): void {
if (ownerId !== undefined) {
if (session.hasFallbackOwner) {
session.ownerIds.delete(sessionId);
session.hasFallbackOwner = false;
}
session.ownerIds.add(ownerId);
return;
}
if (session.hasFallbackOwner || session.ownerIds.size === 0) {
session.ownerIds.add(sessionId);
session.hasFallbackOwner = true;
}
}
@@ -278,31 +295,39 @@ async function acquireSession(
const inFlight = startingSessions.get(sessionKey);
if (inFlight) {
const session = await waitForPromiseWithCancellation(inFlight, options);
attachOwner(session, sessionId, options.kernelOwnerId);
return session;
attachOwner(inFlight, sessionId, options.kernelOwnerId);
return await waitForPromiseWithCancellation(inFlight.promise, options);
}
let startingSession!: StartingJuliaSession;
const startPromise = (async () => {
try {
const kernel = await startKernel(cwd, options);
const session: JuliaSession = {
sessionKey,
sessionId,
kernel,
owners: new Set<string>(),
};
const kernel = await startKernel(cwd, options);
const session: JuliaSession = {
sessionKey,
sessionId,
cwd,
kernel,
ownerIds: new Set(startingSession.ownerIds),
hasFallbackOwner: startingSession.hasFallbackOwner,
};
if (startingSessions.get(sessionKey) === startingSession) {
sessions.set(sessionKey, session);
return session;
} finally {
startingSessions.delete(sessionKey);
}
return session;
})();
startingSessions.set(sessionKey, startPromise);
const session = await waitForPromiseWithCancellation(startPromise, options);
attachOwner(session, sessionId, options.kernelOwnerId);
return session;
startingSession = {
ownerIds: new Set(),
hasFallbackOwner: false,
promise: startPromise,
};
attachOwner(startingSession, sessionId, options.kernelOwnerId);
startingSessions.set(sessionKey, startingSession);
try {
return await waitForPromiseWithCancellation(startPromise, options);
} finally {
if (startingSessions.get(sessionKey) === startingSession) startingSessions.delete(sessionKey);
}
}
async function replaceSessionKernel(session: JuliaSession, cwd: string, options: JuliaExecutorOptions): Promise<void> {
@@ -310,38 +335,107 @@ async function replaceSessionKernel(session: JuliaSession, cwd: string, options:
sessionKey: session.sessionKey,
});
const oldKernel = session.kernel;
void oldKernel.shutdown({ timeoutMs: SHUTDOWN_GRACE_MS }).catch(() => {});
const kernelPromise = startKernel(cwd, options);
const kernel = await waitForPromiseWithCancellation(kernelPromise, options);
session.kernel = kernel;
const remaining = getRemainingTimeoutMs(options.deadlineMs);
await oldKernel
.shutdown(remaining !== undefined ? { timeoutMs: Math.max(0, remaining) } : undefined)
.catch(() => undefined);
if (sessions.get(session.sessionKey) !== session) {
throw new JuliaExecutionCancelledError(false);
}
requireRemainingTimeoutMs(options.deadlineMs);
const nextKernel = await startKernel(cwd, options);
if (sessions.get(session.sessionKey) !== session) {
await nextKernel.shutdown().catch(() => undefined);
throw new JuliaExecutionCancelledError(false);
}
session.kernel = nextKernel;
}
async function resetSession(sessionKey: string): Promise<void> {
const session = sessions.get(sessionKey);
const session = sessions.get(sessionKey) ?? (await startingSessions.get(sessionKey)?.promise.catch(() => undefined));
if (!session) return;
sessions.delete(sessionKey);
await session.kernel.shutdown({ timeoutMs: SHUTDOWN_GRACE_MS }).catch(() => {});
await session.kernel.shutdown({ timeoutMs: SHUTDOWN_GRACE_MS }).catch(() => undefined);
}
export async function disposeAllJuliaKernelSessions(): Promise<void> {
const active = Array.from(sessions.values());
sessions.clear();
const pending = [...startingSessions.values()].map(starting => starting.promise);
startingSessions.clear();
resettingSessions.clear();
await Promise.all(active.map(s => s.kernel.shutdown({ timeoutMs: SHUTDOWN_GRACE_MS }).catch(() => {})));
const started = await Promise.allSettled(pending);
const all = [...sessions.entries()];
for (const result of started) {
if (result.status !== "fulfilled") continue;
if (!all.some(([, session]) => session === result.value)) {
all.push([result.value.sessionKey, result.value]);
}
}
for (const [id, session] of all) {
if (sessions.get(id) === session) sessions.delete(id);
}
const results = await Promise.allSettled(all.map(([, session]) => session.kernel.shutdown()));
for (let i = 0; i < all.length; i += 1) {
const [id, session] = all[i];
const result = results[i];
if (result.status === "fulfilled" && result.value?.confirmed !== false) continue;
const reason = result.status === "rejected" ? result.reason : "not confirmed";
logger.warn("Julia kernel shutdown not confirmed", {
sessionId: session.sessionId,
sessionKey: id,
cwd: session.cwd,
reason,
});
if (!sessions.has(id)) sessions.set(id, session);
}
}
export async function disposeJuliaKernelSessionsByOwner(ownerId: string): Promise<void> {
const victims: JuliaSession[] = [];
for (const [key, session] of sessions) {
session.owners.delete(ownerId);
if (session.owners.size === 0) {
sessions.delete(key);
victims.push(session);
const toShutdown: JuliaSession[] = [];
const startingToShutdown: StartingJuliaSession[] = [];
for (const session of [...sessions.values()]) {
if (!session.ownerIds.has(ownerId)) continue;
if (session.ownerIds.size === 1) {
toShutdown.push(session);
continue;
}
session.ownerIds.delete(ownerId);
}
for (const [sessionKey, starting] of [...startingSessions.entries()]) {
if (sessions.has(sessionKey) || !starting.ownerIds.has(ownerId)) continue;
if (starting.ownerIds.size === 1) {
startingSessions.delete(sessionKey);
startingToShutdown.push(starting);
continue;
}
starting.ownerIds.delete(ownerId);
}
for (const session of toShutdown) {
if (sessions.get(session.sessionKey) === session) sessions.delete(session.sessionKey);
}
const started = await Promise.allSettled(startingToShutdown.map(starting => starting.promise));
for (const result of started) {
if (result.status !== "fulfilled") continue;
const session = result.value;
if (sessions.get(session.sessionKey) === session) sessions.delete(session.sessionKey);
toShutdown.push(session);
}
const results = await Promise.allSettled(toShutdown.map(session => session.kernel.shutdown()));
for (let i = 0; i < toShutdown.length; i += 1) {
const session = toShutdown[i];
const result = results[i];
if (result.status === "fulfilled" && result.value?.confirmed !== false) {
session.ownerIds.delete(ownerId);
continue;
}
const reason = result.status === "rejected" ? result.reason : "not confirmed";
logger.warn("Julia kernel shutdown not confirmed", {
sessionId: session.sessionId,
sessionKey: session.sessionKey,
cwd: session.cwd,
reason,
});
if (!sessions.has(session.sessionKey)) sessions.set(session.sessionKey, session);
}
await Promise.all(victims.map(s => s.kernel.shutdown({ timeoutMs: SHUTDOWN_GRACE_MS }).catch(() => {})));
}
async function executeWithKernel(
+472 -31
View File
@@ -76,6 +76,19 @@ function display_image(base64_str::String, mime_type::String = "image/png")
return nothing
end
function __omp_emit_status(op::String, fields::AbstractDict=Dict{String, Any}())
status = Dict{String, Any}("op" => op)
for (k, v) in fields
status[string(k)] = v
end
Main.emit_frame(Dict(
"type" => "display",
"id" => Main.current_rid,
"bundle" => Dict("application/x-omp-status" => status)
))
return nothing
end
# -------------------------------------------------------------------------
# File helpers
# -------------------------------------------------------------------------
@@ -154,8 +167,9 @@ function append(path, content)
return resolved
end
function tree(path=".", max_depth=3, show_hidden=false)
function tree(path=".", positional_max_depth=3, positional_show_hidden=false; max_depth=positional_max_depth, show_hidden=positional_show_hidden)
base = string(path)
resolved = __omp_resolve_path(base)
lines = String[]
function walk(dir, prefix, depth)
@@ -183,7 +197,7 @@ function tree(path=".", max_depth=3, show_hidden=false)
end
end
walk(base, "", 1)
walk(resolved, "", 1)
out = join(lines, '\n')
Main.emit_frame(Dict(
@@ -192,7 +206,7 @@ function tree(path=".", max_depth=3, show_hidden=false)
"bundle" => Dict(
"application/x-omp-status" => Dict(
"op" => "tree",
"path" => base,
"path" => resolved,
"lines" => length(lines)
)
)
@@ -200,6 +214,366 @@ function tree(path=".", max_depth=3, show_hidden=false)
return out
end
function __omp_lines_keepends(content::String)
parts = split(content, '\n'; keepempty=true)
if length(parts) == 1 && isempty(parts[1])
return String[]
end
lines = String[]
for i in eachindex(parts)
if i < length(parts)
push!(lines, string(parts[i], "\n"))
elseif !isempty(parts[i])
push!(lines, string(parts[i]))
end
end
return lines
end
function __omp_diff_ops(a::Vector{String}, b::Vector{String})
n = length(a)
m = length(b)
ops = Vector{Tuple{Symbol, Int, Int}}()
if n * m > 4_000_000
for i in 1:n
push!(ops, (:delete, i, 1))
end
for j in 1:m
push!(ops, (:insert, n + 1, j))
end
return ops
end
dp = [zeros(Int, m + 1) for _ in 1:(n + 1)]
for i in n:-1:1
for j in m:-1:1
dp[i][j] = a[i] == b[j] ? dp[i + 1][j + 1] + 1 : max(dp[i + 1][j], dp[i][j + 1])
end
end
i = 1
j = 1
while i <= n && j <= m
if a[i] == b[j]
push!(ops, (:equal, i, j))
i += 1
j += 1
elseif dp[i + 1][j] >= dp[i][j + 1]
push!(ops, (:delete, i, j))
i += 1
else
push!(ops, (:insert, i, j))
j += 1
end
end
while i <= n
push!(ops, (:delete, i, j))
i += 1
end
while j <= m
push!(ops, (:insert, i, j))
j += 1
end
return ops
end
function __omp_unified_diff(a::Vector{String}, b::Vector{String}, from_file::String, to_file::String, context::Int=3)
ops = __omp_diff_ops(a, b)
if !any(op -> op[1] != :equal, ops)
return ""
end
entries = [Dict{Symbol, Any}(:tag => tag, :ai => ai, :bi => bi, :text => tag == :insert ? b[bi] : a[ai]) for (tag, ai, bi) in ops]
changed = [i for i in eachindex(entries) if entries[i][:tag] != :equal]
groups = Vector{Tuple{Int, Int}}()
start = nothing
prev = nothing
for idx in changed
if start === nothing
start = idx
prev = idx
elseif idx - prev <= (2 * context) + 1
prev = idx
else
push!(groups, (start, prev))
start = idx
prev = idx
end
end
if start !== nothing
push!(groups, (start, prev))
end
out = IOBuffer()
write(out, "--- $from_file\n")
write(out, "+++ $to_file\n")
for (group_start, group_end) in groups
lo = max(group_start - context, 1)
hi = min(group_end + context, length(entries))
slice = entries[lo:hi]
a_start = nothing
a_count = 0
b_start = nothing
b_count = 0
for entry in slice
if entry[:tag] != :insert
if a_start === nothing
a_start = entry[:ai]
end
a_count += 1
end
if entry[:tag] != :delete
if b_start === nothing
b_start = entry[:bi]
end
b_count += 1
end
end
write(out, "@@ -$(a_start === nothing ? 1 : a_start),$a_count +$(b_start === nothing ? 1 : b_start),$b_count @@\n")
for entry in slice
prefix = entry[:tag] == :equal ? " " : (entry[:tag] == :delete ? "-" : "+")
text = string(entry[:text])
if !endswith(text, "\n")
text *= "\n"
end
write(out, prefix * text)
end
end
return String(take!(out))
end
function Base.diff(a::AbstractString, b::AbstractString)
path_a = __omp_resolve_path(string(a))
path_b = __omp_resolve_path(string(b))
lines_a = __omp_lines_keepends(open(path_a, "r") do io
Base.read(io, String)
end)
lines_b = __omp_lines_keepends(open(path_b, "r") do io
Base.read(io, String)
end)
out = __omp_unified_diff(lines_a, lines_b, path_a, path_b)
__omp_emit_status("diff", Dict{String, Any}(
"file_a" => path_a,
"file_b" => path_b,
"identical" => isempty(out),
"preview" => first(out, min(500, length(out)))
))
return out
end
function __omp_apply_query(data, query)
if query === nothing || isempty(string(query))
return data
end
q = strip(string(query))
if startswith(q, ".")
q = length(q) == 1 ? "" : q[2:end]
end
if isempty(q)
return data
end
tokens = Vector{Tuple{Symbol, Any}}()
buf = ""
chars = collect(q)
i = 1
while i <= length(chars)
ch = chars[i]
if ch == '.'
if !isempty(buf)
push!(tokens, (:key, buf))
buf = ""
end
elseif ch == '['
if !isempty(buf)
push!(tokens, (:key, buf))
buf = ""
end
j = i + 1
while j <= length(chars) && chars[j] != ']'
j += 1
end
inner = j > i + 1 ? String(chars[(i + 1):(j - 1)]) : ""
if startswith(inner, "\"") && endswith(inner, "\"")
push!(tokens, (:key, length(inner) <= 2 ? "" : inner[2:end-1]))
else
push!(tokens, (:index, parse(Int, inner)))
end
i = j
else
buf *= string(ch)
end
i += 1
end
if !isempty(buf)
push!(tokens, (:key, buf))
end
current = data
for (kind, value) in tokens
if kind == :index
if !(current isa AbstractVector)
return nothing
end
idx = Int(value)
julia_idx = idx >= 0 ? idx + 1 : length(current) + idx + 1
if julia_idx < 1 || julia_idx > length(current)
return nothing
end
current = current[julia_idx]
else
key = string(value)
if !(current isa AbstractDict) || !haskey(current, key)
return nothing
end
current = current[key]
end
end
return current
end
function __omp_json_text(value)
if value === nothing
return "null"
end
try
return Main.json_serialize(value)
catch
return string(value)
end
end
function __omp_optional_int(value, default::Int)
if value === nothing
return default
elseif value isa Integer
return Int(value)
elseif value isa AbstractFloat
return Int(trunc(value))
elseif value isa AbstractString
return parse(Int, value)
end
return Int(value)
end
function output(ids...; format="raw", query=nothing, offset=nothing, limit=nothing)
artifacts_dir = get(ENV, "PI_ARTIFACTS_DIR", "")
if isempty(artifacts_dir)
session_file = get(ENV, "PI_SESSION_FILE", "")
if isempty(session_file)
__omp_emit_status("output", Dict{String, Any}("error" => "No session file available"))
error("No session - output artifacts unavailable")
end
artifacts_dir = replace(session_file, r"\.[^.]*$" => "")
end
if !isdir(artifacts_dir)
__omp_emit_status("output", Dict{String, Any}("error" => "Artifacts directory not found", "path" => artifacts_dir))
error("No artifacts directory found: $artifacts_dir")
end
if isempty(ids)
__omp_emit_status("output", Dict{String, Any}("error" => "No IDs provided"))
error("At least one output ID is required")
end
if query !== nothing && (offset !== nothing || limit !== nothing)
__omp_emit_status("output", Dict{String, Any}("error" => "query cannot be combined with offset/limit"))
error("query cannot be combined with offset/limit")
end
results = Vector{Dict{String, Any}}()
not_found = String[]
for output_id_value in ids
output_id = string(output_id_value)
output_path = joinpath(artifacts_dir, output_id * ".md")
if !isfile(output_path)
push!(not_found, output_id)
continue
end
raw = open(output_path, "r") do io
Base.read(io, String)
end
raw_lines = split(raw, '\n'; keepempty=true)
total_lines = length(raw_lines)
selected = raw
range_info = nothing
if query !== nothing
json_value = try
Main.json_parse(raw)
catch err
__omp_emit_status("output", Dict{String, Any}("id" => output_id, "error" => "Not valid JSON: $(err)"))
error("Output $output_id is not valid JSON: $(err)")
end
result_value = __omp_apply_query(json_value, query)
selected = __omp_json_text(result_value)
elseif offset !== nothing || limit !== nothing
start_line = max(1, __omp_optional_int(offset, 1))
if start_line > total_lines
__omp_emit_status("output", Dict{String, Any}("id" => output_id, "error" => "Offset $start_line beyond end ($total_lines lines)"))
error("Offset $start_line is beyond end of output ($total_lines lines) for $output_id")
end
effective_limit = limit === nothing ? total_lines - start_line + 1 : __omp_optional_int(limit, total_lines - start_line + 1)
end_line = min(total_lines, start_line + effective_limit - 1)
selected = join(raw_lines[start_line:end_line], '\n')
range_info = Dict{String, Any}("start_line" => start_line, "end_line" => end_line, "total_lines" => total_lines)
end
if format == "stripped"
selected = replace(selected, r"\x1b\[[0-9;]*m" => "")
end
if format == "json"
entry = Dict{String, Any}(
"id" => output_id,
"path" => output_path,
"line_count" => query !== nothing ? length(split(selected, '\n')) : total_lines,
"char_count" => query !== nothing ? length(selected) : length(raw),
"content" => selected
)
if range_info !== nothing
entry["range"] = range_info
end
if query !== nothing
entry["query"] = query
end
push!(results, entry)
else
push!(results, Dict{String, Any}("id" => output_id, "content" => selected))
end
end
if !isempty(not_found)
available = sort([replace(name, r"\.md$" => "") for name in readdir(artifacts_dir) if endswith(name, ".md")])
msg = "Output not found: $(join(not_found, ", "))"
if !isempty(available)
shown = available[1:min(20, length(available))]
msg *= "\n\nAvailable outputs: $(join(shown, ", "))"
if length(available) > 20
msg *= " (and $(length(available) - 20) more)"
end
end
__omp_emit_status("output", Dict{String, Any}("not_found" => not_found, "available_count" => length(available)))
error(msg)
end
if length(ids) == 1
if format == "json"
__omp_emit_status("output", Dict{String, Any}("id" => string(ids[1]), "chars" => results[1]["char_count"]))
return results[1]
end
__omp_emit_status("output", Dict{String, Any}("id" => string(ids[1]), "chars" => length(results[1]["content"])))
return results[1]["content"]
end
if format == "json"
__omp_emit_status("output", Dict{String, Any}("count" => length(results), "total_chars" => sum(r["char_count"] for r in results)))
return results
end
__omp_emit_status("output", Dict{String, Any}("count" => length(results), "total_chars" => sum(length(r["content"]) for r in results)))
return results
end
function env(key=nothing, value=nothing)
if key === nothing
items = Dict{String, String}()
@@ -344,20 +718,68 @@ const tool = OmpToolProxy()
# Agent calls
# -------------------------------------------------------------------------
function completion(prompt::String; kwargs...)
args_dict = Dict{String, Any}("prompt" => prompt)
function completion(prompt::String; model="default", system=nothing, schema=nothing, kwargs...)
args_dict = Dict{String, Any}("prompt" => prompt, "model" => model)
if system !== nothing
args_dict["system"] = system
end
if schema !== nothing
args_dict["schema"] = schema
end
for (k, v) in kwargs
args_dict[string(k)] = v
end
return __omp_call_bridge("completion", args_dict)
res = __omp_call_bridge("__completion__", args_dict)
text = res isa AbstractDict ? get(res, "text", res) : res
return schema === nothing ? text : Main.json_parse(string(text))
end
function agent(prompt::String; kwargs...)
function agent(prompt::String; agent_type="task", model=nothing, label=nothing, schema=nothing, return_handle=false, kwargs...)
args_dict = Dict{String, Any}("prompt" => prompt)
for (k, v) in kwargs
args_dict[string(k)] = v
if agent_type !== nothing
args_dict["agentType"] = agent_type
end
return __omp_call_bridge("agent", args_dict)
if model !== nothing
args_dict["model"] = model
end
if label !== nothing
args_dict["label"] = label
end
if schema !== nothing
args_dict["schema"] = schema
end
handle_result = return_handle
for (k, v) in kwargs
key = string(k)
if key == "agent_type" || key == "agentType"
args_dict["agentType"] = v
elseif key == "return_handle" || key == "returnHandle"
handle_result = Bool(v)
else
args_dict[key] = v
end
end
res = __omp_call_bridge("__agent__", args_dict)
text = res isa AbstractDict ? get(res, "text", res) : res
parsed = schema === nothing ? text : Main.json_parse(string(text))
if !handle_result
return parsed
end
details = res isa AbstractDict ? get(res, "details", nothing) : nothing
if !(details isa AbstractDict) || get(details, "id", nothing) === nothing
return Dict{String, Any}("text" => text, "output" => text, "handle" => nothing, "id" => nothing, "agent" => nothing)
end
node = Dict{String, Any}(
"text" => text,
"output" => text,
"handle" => "agent://" * string(get(details, "id", nothing)),
"id" => get(details, "id", nothing),
"agent" => get(details, "agent", nothing)
)
if schema !== nothing
node["data"] = parsed
end
return node
end
function Base.log(message::AbstractString)
@@ -394,8 +816,9 @@ end
function _concurrency_limit()
try
limit_val = __omp_call_bridge("concurrency-bridge", Dict{String, Any}())
return limit_val isa Number ? Int(limit_val) : 0
snap = __omp_call_bridge("__concurrency__", Dict{String, Any}())
limit_val = snap isa AbstractDict ? get(snap, "limit", 0) : snap
return limit_val isa Number ? max(Int(limit_val), 0) : 0
catch
return 0
end
@@ -456,34 +879,52 @@ end
# Budget
# -------------------------------------------------------------------------
struct OmpBudgetHardProxy end
struct OmpBudgetProxy end
function Base.getproperty(::OmpBudgetHardProxy, sym::Symbol)
if sym === :total
return __omp_call_bridge("budget:total", Dict{String, Any}())
elseif sym === :spent
return () -> __omp_call_bridge("budget:spent", Dict{String, Any}())
elseif sym === :remaining
return () -> __omp_call_bridge("budget:remaining", Dict{String, Any}())
function __omp_budget_snapshot()
try
snap = __omp_call_bridge("__budget__", Dict{String, Any}())
return snap isa AbstractDict ? snap : Dict{String, Any}()
catch
return Dict{String, Any}()
end
error("Unknown budget hard metric: $sym")
end
struct OmpBudgetProxy
hard::OmpBudgetHardProxy
function __omp_budget_int(value, default::Int=0)
if value isa Integer
return Int(value)
elseif value isa AbstractFloat
return Int(trunc(value))
elseif value isa AbstractString
try
return parse(Int, value)
catch
return default
end
end
return default
end
function Base.getproperty(bp::OmpBudgetProxy, sym::Symbol)
if sym === :hard
return bp.hard
elseif sym === :total
return __omp_call_bridge("budget:total", Dict{String, Any}())
function Base.getproperty(::OmpBudgetProxy, sym::Symbol)
if sym === :total
snap = __omp_budget_snapshot()
return get(snap, "total", nothing)
elseif sym === :hard
snap = __omp_budget_snapshot()
return get(snap, "hard", false) == true
elseif sym === :spent
return () -> __omp_call_bridge("budget:spent", Dict{String, Any}())
return () -> __omp_budget_int(get(__omp_budget_snapshot(), "spent", 0), 0)
elseif sym === :remaining
return () -> __omp_call_bridge("budget:remaining", Dict{String, Any}())
return () -> begin
snap = __omp_budget_snapshot()
total = get(snap, "total", nothing)
if total === nothing
return Inf
end
return max(0, __omp_budget_int(total, 0) - __omp_budget_int(get(snap, "spent", 0), 0))
end
end
error("Unknown budget metric: $sym")
end
const budget = OmpBudgetProxy(OmpBudgetHardProxy())
const budget = OmpBudgetProxy()
+155 -11
View File
@@ -14,6 +14,16 @@ redirect_stdin(devnull)
global current_rid = nothing
const write_lock = ReentrantLock()
const drain_state_lock = ReentrantLock()
mutable struct DrainBarrier
marker::Vector{UInt8}
remaining::Int
done::Channel{Nothing}
end
const active_drain_barrier = Ref{Union{Nothing, DrainBarrier}}(nothing)
const drain_marker_counter = Ref{UInt}(0)
function json_parse(s::String)
chars = collect(s)
@@ -274,17 +284,155 @@ function emit_frame(frame)
end
end
function drain_stream(rd, kind)
function find_subsequence(haystack::Vector{UInt8}, needle::Vector{UInt8})
needle_len = length(needle)
if needle_len == 0 || length(haystack) < needle_len
return nothing
end
last_start = length(haystack) - needle_len + 1
for start in 1:last_start
matched = true
@inbounds for offset in 1:needle_len
if haystack[start + offset - 1] != needle[offset]
matched = false
break
end
end
if matched
return start
end
end
return nothing
end
function marker_overlap(haystack::Vector{UInt8}, needle::Vector{UInt8})
max_overlap = min(length(haystack), length(needle) - 1)
for overlap in max_overlap:-1:1
matched = true
@inbounds for offset in 1:overlap
if haystack[length(haystack) - overlap + offset] != needle[offset]
matched = false
break
end
end
if matched
return overlap
end
end
return 0
end
function emit_stream_bytes(kind, bytes::Vector{UInt8})
rid = current_rid
if rid === nothing || isempty(bytes)
return
end
emit_frame(Dict("type" => kind, "id" => rid, "data" => String(copy(bytes))))
end
function signal_drain_barrier!(marker::Vector{UInt8})
done_channel = lock(drain_state_lock) do
barrier = active_drain_barrier[]
if barrier === nothing || barrier.marker != marker
return nothing
end
barrier.remaining -= 1
if barrier.remaining == 0
active_drain_barrier[] = nothing
return barrier.done
end
return nothing
end
if done_channel !== nothing
put!(done_channel, nothing)
end
return nothing
end
function process_stream_buffer!(buffer::Vector{UInt8}, kind; flush_all::Bool=false)
while true
barrier = lock(drain_state_lock) do
active_drain_barrier[]
end
marker = barrier === nothing ? nothing : barrier.marker
if marker === nothing
if !isempty(buffer)
emit_stream_bytes(kind, buffer)
empty!(buffer)
end
return
end
marker_index = find_subsequence(buffer, marker)
if marker_index !== nothing
emit_len = marker_index - 1
if emit_len > 0
emit_stream_bytes(kind, buffer[1:emit_len])
end
deleteat!(buffer, 1:(marker_index + length(marker) - 1))
signal_drain_barrier!(marker)
continue
end
keep_len = flush_all ? 0 : marker_overlap(buffer, marker)
emit_len = length(buffer) - keep_len
if emit_len > 0
emit_stream_bytes(kind, buffer[1:emit_len])
deleteat!(buffer, 1:emit_len)
end
return
end
end
function next_drain_marker()
drain_marker_counter[] += UInt(1)
return Vector{UInt8}(codeunits("\0__OMP_DRAIN__:" * string(drain_marker_counter[]) * ":" * string(time_ns()) * "\0"))
end
function await_stream_drains()
flush(stdout)
flush(stderr)
barrier = DrainBarrier(next_drain_marker(), 2, Channel{Nothing}(1))
lock(drain_state_lock) do
if active_drain_barrier[] !== nothing
error("Drain barrier already active")
end
active_drain_barrier[] = barrier
end
try
while !eof(rd)
line = readline(rd, keep=true)
rid = current_rid
if rid !== nothing && !isempty(line)
emit_frame(Dict("type" => kind, "id" => rid, "data" => line))
Base.write(out_wr, barrier.marker)
flush(out_wr)
Base.write(err_wr, barrier.marker)
flush(err_wr)
take!(barrier.done)
finally
lock(drain_state_lock) do
if active_drain_barrier[] === barrier
active_drain_barrier[] = nothing
end
end
end
return nothing
end
function drain_stream(rd, kind)
buffer = UInt8[]
try
while true
data = readavailable(rd)
if !isempty(data)
append!(buffer, data)
process_stream_buffer!(buffer, kind)
elseif eof(rd)
break
else
sleep(0.001)
end
end
catch
# ignore
finally
process_stream_buffer!(buffer, kind, flush_all=true)
end
end
@@ -469,11 +617,7 @@ function main()
emit_error(rid, err, catch_backtrace())
end
# Flush stdout and stderr writes before sending done frame
flush(stdout)
flush(stderr)
# Yield to make sure async drain processes all writes
sleep(0.01)
await_stream_drains()
emit_frame(Dict(
"type" => "done",
+51 -20
View File
@@ -98,17 +98,24 @@ export interface RubyResult {
// register against the same tuple; the kernel stays alive until the last owner detaches.
// ---------------------------------------------------------------------------
interface RubySession {
sessionKey: string;
sessionId: string;
cwd: string;
kernel: RubyKernel;
interface RubySessionOwners {
ownerIds: Set<string>;
hasFallbackOwner: boolean;
}
interface RubySession extends RubySessionOwners {
sessionKey: string;
sessionId: string;
cwd: string;
kernel: RubyKernel;
}
interface StartingRubySession extends RubySessionOwners {
promise: Promise<RubySession>;
}
const sessions = new Map<string, RubySession>();
const startingSessions = new Map<string, Promise<RubySession>>();
const startingSessions = new Map<string, StartingRubySession>();
const resettingSessions = new Map<string, Promise<void>>();
function normalizeSessionCwd(cwd: string): string {
@@ -317,7 +324,7 @@ async function startKernel(cwd: string, options: RubyExecutorOptions): Promise<R
});
}
function attachOwner(session: RubySession, sessionId: string, ownerId: string | undefined): void {
function attachOwner(session: RubySessionOwners, sessionId: string, ownerId: string | undefined): void {
if (ownerId !== undefined) {
if (session.hasFallbackOwner) {
session.ownerIds.delete(sessionId);
@@ -345,10 +352,10 @@ async function acquireSession(
}
const starting = startingSessions.get(sessionKey);
if (starting) {
const session = await starting;
attachOwner(session, sessionId, options.kernelOwnerId);
return session;
attachOwner(starting, sessionId, options.kernelOwnerId);
return await starting.promise;
}
let startingSession!: StartingRubySession;
const startup = (async () => {
const kernel = await startKernel(cwd, options);
const session: RubySession = {
@@ -356,19 +363,25 @@ async function acquireSession(
sessionId,
cwd,
kernel,
ownerIds: new Set(),
hasFallbackOwner: false,
ownerIds: new Set(startingSession.ownerIds),
hasFallbackOwner: startingSession.hasFallbackOwner,
};
sessions.set(sessionKey, session);
if (startingSessions.get(sessionKey) === startingSession) {
sessions.set(sessionKey, session);
}
return session;
})();
startingSessions.set(sessionKey, startup);
startingSession = {
ownerIds: new Set(),
hasFallbackOwner: false,
promise: startup,
};
attachOwner(startingSession, sessionId, options.kernelOwnerId);
startingSessions.set(sessionKey, startingSession);
try {
const session = await startup;
attachOwner(session, sessionId, options.kernelOwnerId);
return session;
return await startup;
} finally {
if (startingSessions.get(sessionKey) === startup) startingSessions.delete(sessionKey);
if (startingSessions.get(sessionKey) === startingSession) startingSessions.delete(sessionKey);
}
}
@@ -391,7 +404,8 @@ async function replaceSessionKernel(session: RubySession, cwd: string, options:
}
async function resetSession(sessionKey: string): Promise<void> {
const existing = sessions.get(sessionKey) ?? (await startingSessions.get(sessionKey)?.catch(() => undefined));
const existing =
sessions.get(sessionKey) ?? (await startingSessions.get(sessionKey)?.promise.catch(() => undefined));
if (!existing) return;
sessions.delete(sessionKey);
await existing.kernel.shutdown().catch(() => undefined);
@@ -402,7 +416,7 @@ async function resetSession(sessionKey: string): Promise<void> {
// ---------------------------------------------------------------------------
export async function disposeAllRubyKernelSessions(): Promise<void> {
const pending = [...startingSessions.values()];
const pending = [...startingSessions.values()].map(starting => starting.promise);
startingSessions.clear();
const started = await Promise.allSettled(pending);
const all = [...sessions.entries()];
@@ -433,6 +447,7 @@ export async function disposeAllRubyKernelSessions(): Promise<void> {
export async function disposeRubyKernelSessionsByOwner(ownerId: string): Promise<void> {
const toShutdown: RubySession[] = [];
const startingToShutdown: StartingRubySession[] = [];
for (const session of [...sessions.values()]) {
if (!session.ownerIds.has(ownerId)) continue;
if (session.ownerIds.size === 1) {
@@ -441,9 +456,25 @@ export async function disposeRubyKernelSessionsByOwner(ownerId: string): Promise
}
session.ownerIds.delete(ownerId);
}
for (const [sessionKey, starting] of [...startingSessions.entries()]) {
if (sessions.has(sessionKey) || !starting.ownerIds.has(ownerId)) continue;
if (starting.ownerIds.size === 1) {
startingSessions.delete(sessionKey);
startingToShutdown.push(starting);
continue;
}
starting.ownerIds.delete(ownerId);
}
for (const session of toShutdown) {
if (sessions.get(session.sessionKey) === session) sessions.delete(session.sessionKey);
}
const started = await Promise.allSettled(startingToShutdown.map(starting => starting.promise));
for (const result of started) {
if (result.status !== "fulfilled") continue;
const session = result.value;
if (sessions.get(session.sessionKey) === session) sessions.delete(session.sessionKey);
toShutdown.push(session);
}
const results = await Promise.allSettled(toShutdown.map(session => session.kernel.shutdown()));
for (let i = 0; i < toShutdown.length; i += 1) {
const session = toShutdown[i];
@@ -334,8 +334,13 @@ export class RubyKernel {
});
};
let requestWritten = false;
const requestCancel = () => {
if (pending.settled || pending.escalationTimer) return;
if (!requestWritten) {
finalize();
return;
}
void this.interrupt();
const escalation = setTimeout(() => {
if (pending.settled) return;
@@ -375,6 +380,10 @@ export class RubyKernel {
onAbort();
} else {
options.signal.addEventListener("abort", onAbort, { once: true });
if (options.signal.aborted) {
options.signal.removeEventListener("abort", onAbort);
onAbort();
}
}
}
@@ -389,6 +398,11 @@ export class RubyKernel {
storeHistory: options?.storeHistory ?? !(options?.silent ?? false),
});
if (pending.settled) {
return promise;
}
requestWritten = true;
try {
await this.#writeLine(payload);
} catch (err) {
+33 -19
View File
@@ -16,9 +16,9 @@
# auto-displayed (like IRB) unless it is nil, an assignment, or a definition.
#
# Frame channel isolation: the original stdout is dup'd onto a private IO for
# protocol frames, then fd 1 is repointed at an internal pipe. Child processes
# that inherit fd 1 (any `system`/backtick call) land in that pipe and a drain
# thread re-emits their bytes as stdout frames instead of corrupting the NDJSON
# protocol frames, then fd 1/fd 2 are repointed at internal pipes. Child
# processes that inherit stdout/stderr land in those pipes and drain threads
# re-emit their bytes as stdout/stderr frames instead of corrupting the NDJSON
# channel. Ruby-level writes go through $stdout/$stderr proxies that emit frames
# synchronously so they order correctly with display output.
@@ -39,15 +39,25 @@ $__omp_silent = false
begin
$__omp_frame_io = STDOUT.dup
$__omp_frame_io.sync = true
__omp_cap_r, __omp_cap_w = IO.pipe
STDOUT.reopen(__omp_cap_w)
__omp_stdout_cap_r, __omp_stdout_cap_w = IO.pipe
STDOUT.reopen(__omp_stdout_cap_w)
STDOUT.sync = true
__omp_cap_w.close
$__omp_capture_read = __omp_cap_r
__omp_stdout_cap_w.close
$__omp_stdout_capture_read = __omp_stdout_cap_r
rescue StandardError
$__omp_frame_io = STDOUT
($__omp_frame_io.sync = true) rescue nil
$__omp_capture_read = nil
$__omp_stdout_capture_read = nil
end
begin
__omp_stderr_cap_r, __omp_stderr_cap_w = IO.pipe
STDERR.reopen(__omp_stderr_cap_w)
STDERR.sync = true
__omp_stderr_cap_w.close
$__omp_stderr_capture_read = __omp_stderr_cap_r
rescue StandardError
$__omp_stderr_capture_read = nil
end
# Protect the protocol channel from user code: read requests on a private dup of
@@ -159,8 +169,10 @@ end
# ---------------------------------------------------------------------------
class OmpStreamProxy
def initialize(kind)
def initialize(kind, io, fileno)
@kind = kind
@io = io
@fileno = fileno
end
def write(*args)
@@ -214,19 +226,20 @@ class OmpStreamProxy
def sync=(value); value; end
def tty?; false; end
def isatty; false; end
def fileno; 1; end
def to_io; STDOUT; end
def fileno; @fileno; end
def to_io; @io; end
def closed?; false; end
def fsync; 0; end
def external_encoding; Encoding::UTF_8; end
def external_encoding
(@io.external_encoding rescue Encoding::UTF_8)
end
end
# ---------------------------------------------------------------------------
# fd-1 capture drain (child-process stdout) + parent watchdog
# fd-1/fd-2 capture drains (child-process stdout/stderr) + parent watchdog
# ---------------------------------------------------------------------------
def __omp_start_capture_drain
io = $__omp_capture_read
def __omp_start_capture_drain(io, kind)
return if io.nil?
Thread.new do
loop do
@@ -243,7 +256,7 @@ def __omp_start_capture_drain
if rid.nil?
($__omp_raw_stderr.write(chunk) rescue nil)
else
__omp_emit("type" => "stdout", "id" => rid, "data" => __omp_scrub(chunk))
__omp_emit("type" => kind, "id" => rid, "data" => __omp_scrub(chunk))
end
end
end
@@ -430,11 +443,12 @@ end
# ---------------------------------------------------------------------------
def __omp_main
$stdout = OmpStreamProxy.new("stdout")
$stderr = OmpStreamProxy.new("stderr")
$stdout = OmpStreamProxy.new("stdout", STDOUT, 1)
$stderr = OmpStreamProxy.new("stderr", STDERR, 2)
__omp_install_idle_sigint
__omp_start_parent_watchdog
__omp_start_capture_drain
__omp_start_capture_drain($__omp_stdout_capture_read, "stdout")
__omp_start_capture_drain($__omp_stderr_capture_read, "stderr")
$__omp_proto_stdin.each_line do |raw|
line = raw.strip
+2 -2
View File
@@ -4,13 +4,13 @@ export type EvalLanguage = "python" | "js" | "ruby" | "julia";
import type { ImageContent } from "@oh-my-pi/pi-ai";
import type { OutputMeta } from "../tools/output-meta";
/** Status event emitted by prelude helpers (python or js) for TUI rendering. */
/** Status event emitted by eval prelude helpers for TUI rendering. */
export interface EvalStatusEvent {
op: string;
[key: string]: unknown;
}
/** Display output captured during eval execution. Union of python and js shapes. */
/** Display output captured during eval execution across supported backends. */
export type EvalDisplayOutput =
| { type: "json"; data: unknown }
| { type: "image"; data: string; mimeType: string }
@@ -169,6 +169,7 @@ export type SymbolKey =
| "lang.cpp"
| "lang.csharp"
| "lang.ruby"
| "lang.julia"
| "lang.php"
| "lang.swift"
| "lang.kotlin"
@@ -369,6 +370,7 @@ const UNICODE_SYMBOLS: SymbolMap = {
"lang.cpp": "➕",
"lang.csharp": "♯",
"lang.ruby": "💎",
"lang.julia": "Ⓙ",
"lang.php": "🐘",
"lang.swift": "🕊",
"lang.kotlin": "🅺",
@@ -673,6 +675,7 @@ const NERD_SYMBOLS: SymbolMap = {
"lang.cpp": "\u{E61D}",
"lang.csharp": "\u{E7BC}",
"lang.ruby": "\u{E791}",
"lang.julia": "\u{E624}",
"lang.php": "\u{E608}",
"lang.swift": "\u{E755}",
"lang.kotlin": "\u{E634}",
@@ -870,6 +873,7 @@ const ASCII_SYMBOLS: SymbolMap = {
"lang.cpp": "cpp",
"lang.csharp": "cs",
"lang.ruby": "rb",
"lang.julia": "jl",
"lang.php": "php",
"lang.swift": "swift",
"lang.kotlin": "kt",
@@ -1335,6 +1339,8 @@ const langMap: Record<string, SymbolKey> = {
cs: "lang.csharp",
ruby: "lang.ruby",
rb: "lang.ruby",
julia: "lang.julia",
jl: "lang.julia",
php: "lang.php",
swift: "lang.swift",
kotlin: "lang.kotlin",
@@ -1406,6 +1412,20 @@ const langMap: Record<string, SymbolKey> = {
bin: "lang.binary",
};
/**
* Brand colors for language icons, keyed by the resolved `lang.*` SymbolKey.
* Used by {@link Theme.getLangIconStyled} so eval-kernel cell headers tint each
* language with its recognizable hue (JS yellow, Ruby red, Julia purple, Python
* blue) instead of a flat muted gray. Applied as truecolor/256 per the active
* color mode; languages without an entry fall back to the muted theme color.
*/
const LANG_BRAND_COLORS: Partial<Record<SymbolKey, string>> = {
"lang.javascript": "#f7df1e",
"lang.python": "#3776ab",
"lang.ruby": "#cc342d",
"lang.julia": "#9558b2",
};
/**
* Resolve a theme color value (hex string or 256-color index) to a CSS hex string.
* Empty string represents the default terminal color.
@@ -1874,6 +1894,21 @@ export class Theme {
const key = langMap[normalized];
return key ? this.#symbols[key] : this.#symbols["lang.default"];
}
/**
* Language icon tinted with the language's brand color (see
* {@link LANG_BRAND_COLORS}). Falls back to the muted theme color for
* languages without a brand entry, and returns the bare (possibly empty)
* icon when the active symbol preset has none.
*/
getLangIconStyled(lang: string | undefined): string {
const icon = this.getLangIcon(lang);
if (!icon) return icon;
const key = lang ? langMap[lang.toLowerCase()] : undefined;
const hex = key ? LANG_BRAND_COLORS[key] : undefined;
if (!hex) return this.fg("muted", icon);
return `${colorToAnsi(hex, this.mode)}${icon}\x1b[39m`;
}
}
// ============================================================================
@@ -22,7 +22,7 @@ export type MessageBlock = ({ kind: "code" } & CodeBlock) | ({ kind: "quote" } &
export interface LastCommand {
kind: "bash" | "eval";
code: string;
/** Highlight language: "bash" for bash, "python"/"javascript" for eval. */
/** Highlight language: "bash" for bash, or the resolved eval language ("python"/"javascript"/"ruby"/"julia"). */
language: string;
}
+7 -7
View File
@@ -474,7 +474,7 @@ export interface CreateAgentSessionOptions {
/** Enable LSP integration (tool, formatting, diagnostics, warmup). Default: true */
enableLsp?: boolean;
/** Skip Python kernel availability check and prelude warmup */
/** Skip subprocess-kernel availability checks and prelude warmup */
skipPythonPreflight?: boolean;
/** Tool names explicitly requested (enables disabled-by-default tools) */
toolNames?: string[];
@@ -868,11 +868,11 @@ function registerSshCleanup(): void {
postmortem.register("ssh-cleanup", cleanupSshResources);
}
let pythonCleanupRegistered = false;
let evalCleanupRegistered = false;
function registerPythonCleanup(): void {
if (pythonCleanupRegistered) return;
pythonCleanupRegistered = true;
function registerEvalCleanup(): void {
if (evalCleanupRegistered) return;
evalCleanupRegistered = true;
postmortem.register("python-cleanup", disposeAllKernelSessions);
postmortem.register("ruby-cleanup", disposeAllRubyKernelSessions);
postmortem.register("julia-cleanup", disposeAllJuliaKernelSessions);
@@ -1084,7 +1084,7 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {}
const eventBus = options.eventBus ?? new EventBus();
registerSshCleanup();
registerPythonCleanup();
registerEvalCleanup();
// Pin authStorage to modelRegistry.authStorage: ModelRegistry.getApiKey() routes refresh
// failures through that instance, so any divergent storage handed to the bridge / mcpManager
@@ -2697,7 +2697,7 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {}
const originalDispose = session.dispose.bind(session);
session.dispose = async () => {
try {
// Reject new session work (Python/eval starts) the moment disposal
// Reject new session work (eval starts) the moment disposal
// begins — the lifecycle await below opens an async gap before
// AgentSession.dispose() would otherwise set its guards.
session.beginDispose();
@@ -522,7 +522,7 @@ export interface AgentSessionConfig {
obfuscator?: SecretObfuscator;
/** Inherited eval executor session id from a parent agent. */
parentEvalSessionId?: string;
/** Logical owner for retained Python kernels created by this session. */
/** Logical owner for retained eval kernels created by this session. */
evalKernelOwnerId?: string;
/**
* AsyncJobManager that this session installed as the process-global instance.
@@ -4175,7 +4175,7 @@ export class AgentSession {
/**
* Synchronously mark the session as disposing so new work is rejected
* immediately: Python/eval starts throw, queued asides are dropped, and the
* immediately: eval starts throw, queued asides are dropped, and the
* aside provider is detached. Idempotent; `dispose()` runs it first.
*
* Wrappers that await other teardown before delegating to `dispose()` MUST
@@ -4237,11 +4237,9 @@ export class AgentSession {
AsyncJobManager.setInstance(undefined);
}
}
const pythonExecutionsSettled = await this.#prepareEvalExecutionsForDispose();
if (!pythonExecutionsSettled) {
logger.warn(
"Detaching retained Python kernel ownership during dispose while Python execution is still active",
);
const evalExecutionsSettled = await this.#prepareEvalExecutionsForDispose();
if (!evalExecutionsSettled) {
logger.warn("Detaching retained eval-kernel ownership during dispose while eval execution is still active");
}
await disposeKernelSessionsByOwner(this.#evalKernelOwnerId);
await disposeRubyKernelSessionsByOwner(this.#evalKernelOwnerId);
+4 -6
View File
@@ -41,7 +41,8 @@ import type { AuthStorage } from "../session/auth-storage";
import { SKILL_PROMPT_MESSAGE_TYPE, USER_INTERRUPT_LABEL } from "../session/messages";
import { SessionManager } from "../session/session-manager";
import { truncateTail } from "../session/streaming-output";
import type { ContextFileEntry } from "../tools";
import type { ContextFileEntry, ToolSession } from "../tools";
import { resolveEvalBackends } from "../tools/eval-backends";
import { isIrcEnabled } from "../tools/irc";
import { normalizeSchema } from "../tools/jtd-to-json-schema";
import {
@@ -1790,12 +1791,9 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
toolNames = [...toolNames, "irc"];
}
if (toolNames?.includes("exec")) {
const allowEvalPy = settings.get("eval.py") ?? true;
const allowEvalJs = settings.get("eval.js") ?? true;
const allowEvalRb = settings.get("eval.rb") ?? true;
const allowEvalJl = settings.get("eval.jl") ?? true;
const backends = resolveEvalBackends({ settings } as ToolSession);
const expanded = toolNames.filter(name => name !== "exec");
if (allowEvalPy || allowEvalJs || allowEvalRb || allowEvalJl) expanded.push("eval");
if (backends.python || backends.js || backends.ruby || backends.julia) expanded.push("eval");
expanded.push("bash");
toolNames = Array.from(new Set(expanded));
}
@@ -19,8 +19,9 @@ export function readEvalBackendsAllowance(session: ToolSession): EvalBackendsAll
}
/**
* Materialize the active eval backend allowance: PI_PY / PI_JS / PI_RB env flags
* override the per-key settings; otherwise settings (defaults true) win.
* Materialize the active eval backend allowance: PI_PY / PI_JS / PI_RB / PI_JL
* env flags override the per-key settings; otherwise settings (defaults true)
* win.
*/
export function resolveEvalBackends(session: ToolSession): EvalBackendsAllowance {
const settings = readEvalBackendsAllowance(session);
@@ -2,11 +2,12 @@
* TUI rendering for the eval tool.
*
* Split out from `eval.ts` so the renderer can be imported by `renderers.ts`
* without dragging the eval *runtime* (JS/Python backends -> agent bridge ->
* task executor -> sdk -> extension loader -> root barrel) into the renderer
* module graph. That transitive chain re-enters `renderers.ts` while `eval.ts`
* is still initializing, which previously crashed module load with a TDZ
* `Cannot access 'evalToolRenderer' before initialization`.
* without dragging the eval *runtime* (JS/Python/Ruby/Julia backends ->
* agent bridge -> task executor -> sdk -> extension loader -> root barrel)
* into the renderer module graph. That transitive chain re-enters
* `renderers.ts` while `eval.ts` is still initializing, which previously
* crashed module load with a TDZ `Cannot access 'evalToolRenderer' before
* initialization`.
*/
import type { Component } from "@oh-my-pi/pi-tui";
import { Markdown, Text } from "@oh-my-pi/pi-tui";
@@ -512,6 +513,7 @@ export const evalToolRenderer = {
{
code: cell.code,
language: languageForHighlighter(cell.language),
showLanguage: true,
index: i,
total: cells.length,
title: cell.title,
@@ -614,6 +616,7 @@ export const evalToolRenderer = {
{
code: cell.code,
language: languageForHighlighter(cell.language ?? details?.language),
showLanguage: true,
index: i,
total: cellResults.length,
title: cell.title,
+38 -8
View File
@@ -31,7 +31,9 @@ const evalCellSchema = type({
language: type("'py' | 'js' | 'rb' | 'jl'").describe(
'runtime: "py" for the IPython kernel, "js" for the persistent JS VM, "rb" for the persistent Ruby kernel, "jl" for the persistent Julia kernel',
),
code: type("string").describe("cell body, verbatim. Top-level `await` is available in py/js."),
code: type("string").describe(
"cell body, verbatim. Top-level `await` is available in py/js; rb/jl auto-display the last expression like a REPL.",
),
"title?": type("string").describe('short label shown in transcript (e.g. "imports", "load config")'),
"timeout?": type("number").describe("per-cell timeout in seconds"),
"reset?": type("boolean").describe(
@@ -149,12 +151,18 @@ async function resolveBackend(session: ToolSession, language: EvalLanguage): Pro
const allowPy = backends.python;
const allowJs = backends.js;
const allowRb = backends.ruby;
const allowJl = backends.julia;
if (language === "python") {
if (!allowPy) throw new ToolError("Python backend is disabled (PI_PY=0 or eval.py = false).");
if (!(await pythonBackend.isAvailable(session))) {
const alternatives = [allowJs ? '"js"' : null, allowRb ? '"rb"' : null, allowJl ? '"jl"' : null].filter(
Boolean,
);
throw new ToolError(
'Python backend is unavailable in this session. Pass language: "js" or install the python kernel.',
alternatives.length > 0
? `Python backend is unavailable in this session. Pass language: ${alternatives.join(" or ")} or install the python kernel.`
: 'Python backend is unavailable in this session. Install the python kernel to use language: "py".',
);
}
return { backend: pythonBackend };
@@ -162,20 +170,41 @@ async function resolveBackend(session: ToolSession, language: EvalLanguage): Pro
if (language === "ruby") {
if (!allowRb) throw new ToolError("Ruby backend is disabled (PI_RB=0 or eval.rb = false).");
if (!(await rubyBackend.isAvailable(session))) {
throw new ToolError('Ruby backend is unavailable in this session. Pass language: "js" or install Ruby.');
const alternatives = [allowJs ? '"js"' : null, allowPy ? '"py"' : null, allowJl ? '"jl"' : null].filter(
Boolean,
);
throw new ToolError(
alternatives.length > 0
? `Ruby backend is unavailable in this session. Pass language: ${alternatives.join(" or ")} or install Ruby.`
: 'Ruby backend is unavailable in this session. Install Ruby to use language: "rb".',
);
}
return { backend: rubyBackend };
}
if (language === "julia") {
if (!backends.julia) throw new ToolError("Julia backend is disabled (PI_JL=0 or eval.jl = false).");
if (!allowJl) throw new ToolError("Julia backend is disabled (PI_JL=0 or eval.jl = false).");
if (!(await juliaBackend.isAvailable(session))) {
throw new ToolError('Julia backend is unavailable in this session. Pass language: "js" or install Julia.');
const alternatives = [allowJs ? '"js"' : null, allowPy ? '"py"' : null, allowRb ? '"rb"' : null].filter(
Boolean,
);
throw new ToolError(
alternatives.length > 0
? `Julia backend is unavailable in this session. Pass language: ${alternatives.join(" or ")} or install Julia.`
: 'Julia backend is unavailable in this session. Install Julia to use language: "jl".',
);
}
return { backend: juliaBackend };
}
if (!allowJs) throw new ToolError("JavaScript backend is disabled (PI_JS=0 or eval.js = false).");
return { backend: jsBackend };
}
function formatEvalInputLanguage(value: string): string {
if (value === "py" || value === "python") return "python";
if (value === "js" || value === "javascript") return "javascript";
if (value === "rb" || value === "ruby") return "ruby";
if (value === "jl" || value === "julia") return "julia";
return value;
}
export class EvalTool implements AgentTool<typeof evalSchema> {
readonly name = "eval";
@@ -185,7 +214,8 @@ export class EvalTool implements AgentTool<typeof evalSchema> {
const cells = Array.isArray(params.cells) ? params.cells : [];
const firstCell = cells[0] as Partial<EvalCellInput> | undefined;
if (!firstCell) return [];
const language = typeof firstCell.language === "string" ? firstCell.language : "(missing)";
const language =
typeof firstCell.language === "string" ? formatEvalInputLanguage(firstCell.language) : "javascript (default)";
const code = typeof firstCell.code === "string" ? firstCell.code : "";
const lines = [`Language: ${language}`, `Code:\n${truncateForPrompt(code)}`];
if (cells.length > 1) {
@@ -236,7 +266,7 @@ export class EvalTool implements AgentTool<typeof evalSchema> {
const first = cells.find(c => c && typeof c === "object");
if (!first) return "evaluating";
const title = typeof first.title === "string" ? first.title : undefined;
const language = typeof first.language === "string" ? first.language : "?";
const language = typeof first.language === "string" ? formatEvalInputLanguage(first.language) : "javascript";
const label = title || `running ${language}`;
return cells.length > 1 ? `${label} (+${cells.length - 1})` : label;
};
@@ -382,7 +412,7 @@ export class EvalTool implements AgentTool<typeof evalSchema> {
// The per-cell `timeout` is a budget on the cell runtime's *own*
// work. Host-side `agent()`/`parallel()`/`completion()` bridge calls suspend
// that budget entirely and restart a fresh timeout window when control
// returns to Python/JS. Compute, stdout, `log()`/`phase()`, and
// returns to the active backend runtime. Compute, stdout, `log()`/`phase()`, and
// ordinary tool calls all count against the budget. The watchdog drives
// `combinedSignal`; we pass no wall-clock deadline downstream so the
// backends never arm a competing fixed timer.
+3 -3
View File
@@ -167,7 +167,7 @@ export interface ToolSession {
suppressSpawnAdvisory?: boolean;
/** Optional fetch implementation injected into the URL read pipeline (tests, proxies). Defaults to global fetch. */
fetch?: FetchImpl;
/** Skip Python kernel availability check and warmup */
/** Skip subprocess-kernel availability checks and warmup */
skipPythonPreflight?: boolean;
/** Pre-loaded context files (AGENTS.md, etc) */
contextFiles?: ContextFileEntry[];
@@ -204,13 +204,13 @@ export interface ToolSession {
requireYieldTool?: boolean;
/** Task recursion depth (0 = top-level, 1 = first child, etc.) */
taskDepth?: number;
/** Get shared eval executor session ID. Subagents inherit this to share JS/Python state. */
/** Get shared eval executor session ID. Subagents inherit this to share JS/Python/Ruby/Julia state. */
getEvalSessionId?: () => string | null;
/** Get session file */
getSessionFile: () => string | null;
/** Get eval kernel owner ID for session-scoped retained-kernel cleanup. */
getEvalKernelOwnerId?: () => string | null;
/** Reject new eval (python or js) work once session disposal has started. */
/** Reject new eval work once session disposal has started. */
assertEvalExecutionAllowed?: () => void;
/** Track tool-owned eval work so session disposal can await/abort it like direct session eval runs. */
trackEvalExecution?<T>(execution: Promise<T>, abortController: AbortController): Promise<T>;
+12 -1
View File
@@ -32,6 +32,13 @@ export interface CodeCellOptions {
*/
codeTail?: boolean;
expanded?: boolean;
/**
* Prefix the header with the cell's language icon (resolved through the
* active symbol preset: nerd-font devicon, unicode emoji, or ascii
* shorthand). Opt-in so only the eval kernel renderer labels each cell;
* read/write/browser code cells stay icon-free.
*/
showLanguage?: boolean;
width: number;
codeStartLine?: number;
codeLineNumbers?: Array<number | null>;
@@ -47,8 +54,12 @@ function getState(status?: CodeCellOptions["status"]): State | undefined {
}
function formatHeader(options: CodeCellOptions, theme: Theme): { title: string; meta?: string } {
const { index, total, title, status, spinnerFrame, duration } = options;
const { index, total, title, status, spinnerFrame, duration, language, showLanguage } = options;
const parts: string[] = [];
if (showLanguage && language) {
const langIcon = theme.getLangIconStyled(language);
if (langIcon) parts.push(langIcon);
}
if (status) {
const icon = formatStatusIcon(
status === "complete"
@@ -85,6 +85,29 @@ describe.skipIf(!SHOULD_RUN)("ruby runner subprocess", () => {
}
});
it("does not write an already-aborted request and keeps the kernel usable", async () => {
using tempDir = TempDir.createSync("@ruby-runner-preabort-");
const kernel = await RubyKernel.start({ cwd: tempDir.path() });
try {
const controller = new AbortController();
controller.abort();
const cancelled = await executeRubyWithKernel(kernel, 'File.write("aborted.txt", "nope")', {
signal: controller.signal,
});
expect(cancelled.cancelled).toBe(true);
const sideEffect = await executeRubyWithKernel(kernel, 'File.exist?("aborted.txt")', {});
expect(sideEffect.exitCode).toBe(0);
expect(sideEffect.output).toContain("false");
const after = await executeRubyWithKernel(kernel, "20 + 22", {});
expect(after.exitCode).toBe(0);
expect(after.output).toContain("42");
} finally {
await kernel.shutdown();
}
});
it("surfaces Ruby errors with a non-zero exit and the message", async () => {
using tempDir = TempDir.createSync("@ruby-runner-error-");
const kernel = await RubyKernel.start({ cwd: tempDir.path() });
@@ -97,6 +120,30 @@ describe.skipIf(!SHOULD_RUN)("ruby runner subprocess", () => {
}
});
it("surfaces Ruby and child-process stderr in output", async () => {
using tempDir = TempDir.createSync("@ruby-runner-stderr-");
const kernel = await RubyKernel.start({ cwd: tempDir.path() });
try {
const result = await executeRubyWithKernel(
kernel,
[
'$stderr.print("ruby stderr\\n")',
'STDERR.print("constant stderr\\n")',
'require "rbconfig"',
'system(RbConfig.ruby, "-e", \'STDERR.print("child stderr\\\\n"); STDOUT.print("child stdout\\\\n")\')',
].join("\n"),
{},
);
expect(result.exitCode).toBe(0);
expect(result.output).toContain("ruby stderr\n");
expect(result.output).toContain("constant stderr\n");
expect(result.output).toContain("child stderr\n");
expect(result.output).toContain("child stdout\n");
} finally {
await kernel.shutdown();
}
});
it("exposes prelude file + text helpers", async () => {
using tempDir = TempDir.createSync("@ruby-runner-prelude-");
const kernel = await RubyKernel.start({ cwd: tempDir.path() });