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:
@@ -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");
|
||||
});
|
||||
});
|
||||
@@ -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(
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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>;
|
||||
|
||||
@@ -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() });
|
||||
|
||||
Reference in New Issue
Block a user