diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index feeee6d42..24a740924 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -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.` 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.` bridge proxy, and bridge-backed `completion`/`agent`/`parallel`/`pipeline`/`log`/`phase`/`budget`, including structured `schema` parsing and `return_handle` DAG shaping). ### Changed diff --git a/packages/coding-agent/src/config/settings-schema.ts b/packages/coding-agent/src/config/settings-schema.ts index 628b6a1a1..8f8ce0772 100644 --- a/packages/coding-agent/src/config/settings-schema.ts +++ b/packages/coding-agent/src/config/settings-schema.ts @@ -122,7 +122,7 @@ export const TAB_GROUPS: Record = { 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.", diff --git a/packages/coding-agent/src/eval/__tests__/julia-prelude.test.ts b/packages/coding-agent/src/eval/__tests__/julia-prelude.test.ts new file mode 100644 index 000000000..fca3f0bfe --- /dev/null +++ b/packages/coding-agent/src/eval/__tests__/julia-prelude.test.ts @@ -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"); + }); +}); diff --git a/packages/coding-agent/src/eval/jl/executor.ts b/packages/coding-agent/src/eval/jl/executor.ts index db72a7e6e..dac5bb6e5 100644 --- a/packages/coding-agent/src/eval/jl/executor.ts +++ b/packages/coding-agent/src/eval/jl/executor.ts @@ -56,11 +56,20 @@ export interface JuliaResult { stdinRequested: boolean; } -interface JuliaSession { +interface JuliaSessionOwners { + ownerIds: Set; + hasFallbackOwner: boolean; +} + +interface JuliaSession extends JuliaSessionOwners { sessionKey: string; sessionId: string; + cwd: string; kernel: JuliaKernel; - owners: Set; +} + +interface StartingJuliaSession extends JuliaSessionOwners { + promise: Promise; } class JuliaExecutionCancelledError extends Error { @@ -71,7 +80,7 @@ class JuliaExecutionCancelledError extends Error { } const sessions = new Map(); -const startingSessions = new Map>(); +const startingSessions = new Map(); const resettingSessions = new Map>(); function normalizeSessionCwd(cwd: string): string { @@ -239,6 +248,7 @@ function buildKernelEnv(options: { } async function startKernel(cwd: string, options: JuliaExecutorOptions): Promise { + requireRemainingTimeoutMs(options.deadlineMs); const env: Record = {}; 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(), - }; + 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 { @@ -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 { - 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 { - 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 { - 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( diff --git a/packages/coding-agent/src/eval/jl/prelude.jl b/packages/coding-agent/src/eval/jl/prelude.jl index 698ba3b99..52ff323ef 100644 --- a/packages/coding-agent/src/eval/jl/prelude.jl +++ b/packages/coding-agent/src/eval/jl/prelude.jl @@ -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() diff --git a/packages/coding-agent/src/eval/jl/runner.jl b/packages/coding-agent/src/eval/jl/runner.jl index d63c4949d..62c466e51 100644 --- a/packages/coding-agent/src/eval/jl/runner.jl +++ b/packages/coding-agent/src/eval/jl/runner.jl @@ -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", diff --git a/packages/coding-agent/src/eval/rb/executor.ts b/packages/coding-agent/src/eval/rb/executor.ts index aeedbf966..cfaceeef0 100644 --- a/packages/coding-agent/src/eval/rb/executor.ts +++ b/packages/coding-agent/src/eval/rb/executor.ts @@ -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; hasFallbackOwner: boolean; } +interface RubySession extends RubySessionOwners { + sessionKey: string; + sessionId: string; + cwd: string; + kernel: RubyKernel; +} + +interface StartingRubySession extends RubySessionOwners { + promise: Promise; +} + const sessions = new Map(); -const startingSessions = new Map>(); +const startingSessions = new Map(); const resettingSessions = new Map>(); function normalizeSessionCwd(cwd: string): string { @@ -317,7 +324,7 @@ async function startKernel(cwd: string, options: RubyExecutorOptions): Promise { 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 { - 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 { // --------------------------------------------------------------------------- export async function disposeAllRubyKernelSessions(): Promise { - 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 { export async function disposeRubyKernelSessionsByOwner(ownerId: string): Promise { 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]; diff --git a/packages/coding-agent/src/eval/rb/kernel.ts b/packages/coding-agent/src/eval/rb/kernel.ts index b00b0a472..ed0664148 100644 --- a/packages/coding-agent/src/eval/rb/kernel.ts +++ b/packages/coding-agent/src/eval/rb/kernel.ts @@ -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) { diff --git a/packages/coding-agent/src/eval/rb/runner.rb b/packages/coding-agent/src/eval/rb/runner.rb index 9fa80d616..cc497d292 100644 --- a/packages/coding-agent/src/eval/rb/runner.rb +++ b/packages/coding-agent/src/eval/rb/runner.rb @@ -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 diff --git a/packages/coding-agent/src/eval/types.ts b/packages/coding-agent/src/eval/types.ts index df22ccf0b..2463ada07 100644 --- a/packages/coding-agent/src/eval/types.ts +++ b/packages/coding-agent/src/eval/types.ts @@ -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 } diff --git a/packages/coding-agent/src/modes/theme/theme.ts b/packages/coding-agent/src/modes/theme/theme.ts index 9b39345e9..ca5879223 100644 --- a/packages/coding-agent/src/modes/theme/theme.ts +++ b/packages/coding-agent/src/modes/theme/theme.ts @@ -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 = { 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 = { 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> = { + "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`; + } } // ============================================================================ diff --git a/packages/coding-agent/src/modes/utils/copy-targets.ts b/packages/coding-agent/src/modes/utils/copy-targets.ts index b93c9129f..3eb0bdce0 100644 --- a/packages/coding-agent/src/modes/utils/copy-targets.ts +++ b/packages/coding-agent/src/modes/utils/copy-targets.ts @@ -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; } diff --git a/packages/coding-agent/src/sdk.ts b/packages/coding-agent/src/sdk.ts index acc71956f..8b57ac801 100644 --- a/packages/coding-agent/src/sdk.ts +++ b/packages/coding-agent/src/sdk.ts @@ -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(); diff --git a/packages/coding-agent/src/session/agent-session.ts b/packages/coding-agent/src/session/agent-session.ts index a2061f4b8..b25b30507 100644 --- a/packages/coding-agent/src/session/agent-session.ts +++ b/packages/coding-agent/src/session/agent-session.ts @@ -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); diff --git a/packages/coding-agent/src/task/executor.ts b/packages/coding-agent/src/task/executor.ts index cd10f7433..6837f5d10 100644 --- a/packages/coding-agent/src/task/executor.ts +++ b/packages/coding-agent/src/task/executor.ts @@ -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 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)); } diff --git a/packages/coding-agent/src/tools/eval-backends.ts b/packages/coding-agent/src/tools/eval-backends.ts index 81f808a2f..5888eee6e 100644 --- a/packages/coding-agent/src/tools/eval-backends.ts +++ b/packages/coding-agent/src/tools/eval-backends.ts @@ -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); diff --git a/packages/coding-agent/src/tools/eval-render.ts b/packages/coding-agent/src/tools/eval-render.ts index 98221e503..293b2cd7f 100644 --- a/packages/coding-agent/src/tools/eval-render.ts +++ b/packages/coding-agent/src/tools/eval-render.ts @@ -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, diff --git a/packages/coding-agent/src/tools/eval.ts b/packages/coding-agent/src/tools/eval.ts index 75acb19af..d161c718c 100644 --- a/packages/coding-agent/src/tools/eval.ts +++ b/packages/coding-agent/src/tools/eval.ts @@ -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 { readonly name = "eval"; @@ -185,7 +214,8 @@ export class EvalTool implements AgentTool { const cells = Array.isArray(params.cells) ? params.cells : []; const firstCell = cells[0] as Partial | 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 { 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 { // 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. diff --git a/packages/coding-agent/src/tools/index.ts b/packages/coding-agent/src/tools/index.ts index 9f83f9ae5..2eb00921e 100644 --- a/packages/coding-agent/src/tools/index.ts +++ b/packages/coding-agent/src/tools/index.ts @@ -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?(execution: Promise, abortController: AbortController): Promise; diff --git a/packages/coding-agent/src/tui/code-cell.ts b/packages/coding-agent/src/tui/code-cell.ts index e269a9276..143e78f3d 100644 --- a/packages/coding-agent/src/tui/code-cell.ts +++ b/packages/coding-agent/src/tui/code-cell.ts @@ -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; @@ -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" diff --git a/packages/coding-agent/test/core/ruby-runner.integration.test.ts b/packages/coding-agent/test/core/ruby-runner.integration.test.ts index ea6b29741..c7ccd1fa4 100644 --- a/packages/coding-agent/test/core/ruby-runner.integration.test.ts +++ b/packages/coding-agent/test/core/ruby-runner.integration.test.ts @@ -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() });