From 93e1ce655508cb23f7520001207858e5d05ecae1 Mon Sep 17 00:00:00 2001 From: "prompt.ac/@jeffrey" Date: Thu, 1 Oct 2026 15:30:39 -0700 Subject: [PATCH] Reconnect Codex status colors to live shared-daemon sessions --- slab/bin/codex-session-watch.mjs | 112 +++++++++++++++--- slab/bin/slab-repair-codex-watchers.mjs | 32 +++++ .../Sources/SlabMenubar/ClaudeSession.swift | 4 +- slab/menubar-swift/deploy-host.sh | 1 + slab/menubar-swift/install.sh | 5 +- slab/test/codex-session-watch.test.mjs | 36 ++++++ 6 files changed, 169 insertions(+), 21 deletions(-) create mode 100755 slab/bin/slab-repair-codex-watchers.mjs create mode 100644 slab/test/codex-session-watch.test.mjs diff --git a/slab/bin/codex-session-watch.mjs b/slab/bin/codex-session-watch.mjs index 0f8cf29e78..cc28343b39 100755 --- a/slab/bin/codex-session-watch.mjs +++ b/slab/bin/codex-session-watch.mjs @@ -15,11 +15,13 @@ import { readFile, writeFile, stat, unlink, utimes } from "node:fs/promises"; import { realpathSync } from "node:fs"; import { execFile } from "node:child_process"; import { promisify } from "node:util"; -import { join } from "node:path"; +import { join, basename, resolve } from "node:path"; +import { fileURLToPath } from "node:url"; import { homedir } from "node:os"; -const [sid, _beginArg, wrapperArg, tty = "", cwd = ""] = process.argv.slice(2); -if (!sid) process.exit(1); +const [sid = "", _beginArg, wrapperArg, tty = "", cwd = ""] = process.argv.slice(2); +const isMain = process.argv[1] && realpathSync(process.argv[1]) === fileURLToPath(import.meta.url); +if (isMain && !sid) process.exit(1); const wrapperPid = Number(wrapperArg) || 0; const SLAB_HOME = process.env.SLAB_HOME || join(homedir(), ".local", "share", "slab"); @@ -28,6 +30,7 @@ const AWAITING = join(SLAB_HOME, "state", "awaiting-prompts", sid); const RUNNING = join(SLAB_HOME, "state", "running-tools", sid); const OPEN_IMAGES = join(SLAB_HOME, "state", "open-images"); const SESSIONS = join(process.env.CODEX_HOME || join(homedir(), ".codex"), "sessions"); +const CODEX_ROOT = process.env.CODEX_HOME || join(homedir(), ".codex"); const LOOPBOY_CONFIG = join(homedir(), ".config", "slab", "loopboy.json"); const execFileAsync = promisify(execFile); @@ -47,8 +50,54 @@ const wrapperAlive = () => { // cwd and start within the same second, so "newest file" cross-wires their // rocks. The wrapper and Codex are parent/child; Codex keeps its own rollout // open for writing, giving us an exact, resume-safe association via lsof. +export function rolloutFromTitle(rows, title, cwd) { + const matches = rows.filter(row => typeof row.name === "string" && row.name.length >= 4 + && resolve(row.cwd) === resolve(cwd) + && (title === row.name || title.endsWith(`${row.name} | ${basename(row.cwd)}`))); + // Never guess between concurrent windows or similarly named threads. + return matches.length === 1 ? matches[0].rollout_path : null; +} + +async function namedRollout() { + if (!/^ttys[0-9]+$/.test(tty)) return null; + try { + // New Codex versions share a daemon: its open files belong to MANY TUIs. + // Match this TTY's exact displayed thread name to local thread metadata. + // Read names/paths only; no prompts, transcript text, or credentials. + const { stdout: title } = await execFileAsync("/usr/bin/osascript", ["-e", ` + if application "Terminal" is not running then return "" + tell application "Terminal" + repeat with w in windows + if tty of selected tab of w is "/dev/${tty}" then return name of w + end repeat + end tell + return ""`], { timeout: 2000 }); + if (!title.trim()) return null; + const { stdout } = await execFileAsync("/usr/bin/sqlite3", ["-readonly", "-json", + join(CODEX_ROOT, "state_5.sqlite"), + "select name,cwd,rollout_path from threads where archived=0 and name is not null"], { timeout: 2000 }); + const path = rolloutFromTitle(JSON.parse(stdout || "[]"), title.trim(), cwd); + return path?.startsWith(SESSIONS + "/") ? path : null; + } catch { return null; } +} + async function findRollout() { - for (let i = 0; i < 40 && wrapperAlive(); i++) { + // Reinstalling Slab reattaches watchers to the same live wrappers. Once an + // exact identity is known, prefer it even while Slab owns the window title. + try { + const marker = JSON.parse(await readFile(ACTIVE, "utf8")); + const id = marker.provider_session_id || marker.codex_session_id; + if (/^[0-9a-f]{8}-(?:[0-9a-f]{4}-){3}[0-9a-f]{12}$/i.test(id || "")) { + const { stdout } = await execFileAsync("/usr/bin/sqlite3", ["-readonly", "-json", + join(CODEX_ROOT, "state_5.sqlite"), `select rollout_path from threads where id='${id}'`], { timeout: 2000 }); + const path = JSON.parse(stdout || "[]")[0]?.rollout_path; + if (path?.startsWith(SESSIONS + "/") && (await stat(path)).isFile()) return path; + } + } catch {} + // A resume picker can remain open indefinitely; keep waiting for selection. + for (let i = 0; wrapperAlive(); i++) { + const named = await namedRollout(); + if (named) return named; try { // BSD/macOS ps has no Linux `-P ` selector. Read its compact // PID/PPID table and select the wrapper's children ourselves. @@ -75,19 +124,22 @@ async function findRollout() { const pids = [...descendants] .filter((pid) => pid !== wrapperPid && pid !== process.pid) .map(String); + const rollouts = new Set(); for (const pid of pids) { - const { stdout } = await execFileAsync("/usr/sbin/lsof", ["-Fn", "-p", pid]); - const rollout = stdout.split("\n") + let stdout; + try { ({ stdout } = await execFileAsync("/usr/sbin/lsof", ["-Fn", "-p", pid], { timeout: 2000 })); } + catch { continue; } + for (const rollout of stdout.split("\n") .filter((line) => line.startsWith("n")) .map((line) => line.slice(1)) - .find((p) => p.startsWith(SESSIONS + "/") - && p.includes("/rollout-") && p.endsWith(".jsonl")); - if (rollout) return rollout; + .filter((p) => p.startsWith(SESSIONS + "/") + && p.includes("/rollout-") && p.endsWith(".jsonl"))) rollouts.add(rollout); } + if (rollouts.size === 1) return [...rollouts][0]; } catch { // Codex may not have opened its rollout yet; retry below. } - await sleep(500); + await sleep(i < 20 ? 500 : 5000); } return null; } @@ -182,7 +234,7 @@ async function onVisualArtifact(path) { await onAwaiting("visual artifact ready for review"); } -function handleLine(line, ctx) { +export function handleLine(line, ctx) { // Image generation returns a saved host path beside an inline bitmap. Match // only tool response items so quoted history cannot reopen stale artifacts. const isImageToolOutput = line.includes('"type":"response_item"') @@ -190,7 +242,7 @@ function handleLine(line, ctx) { const generated = isImageToolOutput ? line.match(/Generated images are saved[^\n]*? as (\/[^\s]+\.(?:png|jpe?g|webp))/i) : null; - if (generated) ctx.pending.push(() => onVisualArtifact(generated[1])); + if (generated && !ctx.replaying) ctx.pending.push(() => onVisualArtifact(generated[1])); let obj; try { obj = JSON.parse(line); } catch { return; } const type = obj.type; @@ -204,7 +256,7 @@ function handleLine(line, ctx) { // simply "codex" for the entire first turn. Update the visible subject // when the actual user message arrives; later user items in the same // rollout naturally win over injected context blocks. - ctx.pending.push(() => updateMarker({ + if (!ctx.replaying) ctx.pending.push(() => updateMarker({ subject: t.slice(0, 140), summary: summarize(t), })); @@ -212,6 +264,7 @@ function handleLine(line, ctx) { return; } if (type === "response_item" && payload.role === "assistant") { + if (ctx.replaying) return; const t = textOf(payload); if (/\bRESPONDING\b/i.test(t)) { ctx.pending.push(() => updateMarker({ @@ -227,11 +280,29 @@ function handleLine(line, ctx) { } if (type === "event_msg") { const pt = payload.type || ""; - if (pt === "task_started" || pt === "user_turn") ctx.pending.push(() => onTurnStart(ctx.lastUser)); + const transition = fn => { + if (ctx.replaying) ctx.pending = [fn]; + else ctx.pending.push(fn); + }; + if (pt === "task_started" || pt === "user_turn") { + ctx.turnActive = true; + transition(() => onTurnStart(ctx.lastUser)); + } else if (pt === "task_complete" || pt === "turn_complete") { - ctx.pending.push(() => onTurnComplete(payload.last_agent_message || "")); + ctx.turnActive = false; + transition(() => onTurnComplete(payload.last_agent_message || "")); + } + else if (pt === "turn_aborted") { + ctx.turnActive = false; + transition(async () => { + await rm(RUNNING); await rm(AWAITING); + await updateMarker({ state: "interrupted" }); + }); + } + else if (pt.includes("approval") || pt.includes("elicitation")) { + ctx.turnActive = false; + transition(() => onAwaiting("codex needs approval")); } - else if (pt.includes("approval") || pt.includes("elicitation")) ctx.pending.push(() => onAwaiting("codex needs approval")); } } @@ -245,6 +316,7 @@ async function main() { if (providerId) await updateMarker({ codex_session_id: providerId, provider_session_id: providerId, + transcript_path: file, }); } catch {} // Replay once from the beginning so a resumed Codex window immediately @@ -252,7 +324,7 @@ async function main() { // or aging into interrupted until the user submits another prompt. After // that first pass `offset` makes this an ordinary incremental tail. let offset = 0; - const ctx = { lastUser: "", pending: [], turnActive: false, lastHeartbeatAt: 0 }; + const ctx = { lastUser: "", pending: [], turnActive: false, lastHeartbeatAt: 0, replaying: true }; while (wrapperAlive()) { let size = offset; try { size = (await stat(file)).size; } catch { break; } @@ -265,9 +337,13 @@ async function main() { const tail = chunk.subarray(offset).toString("utf8"); offset = chunk.length; for (const line of tail.split("\n")) if (line.trim()) handleLine(line, ctx); + if (ctx.replaying && ctx.lastUser) { + await updateMarker({ subject: ctx.lastUser.slice(0, 140), summary: summarize(ctx.lastUser) }); + } // Apply transitions in order; last one wins the visible state. for (const fn of ctx.pending) await fn(); ctx.pending = []; + ctx.replaying = false; } // Codex can spend many minutes inside one tool call without appending a // new rollout event. Keep both marker channels fresh for the full active @@ -283,7 +359,7 @@ async function main() { } } -main().catch((error) => { +if (isMain) main().catch((error) => { console.error(`codex-session-watch: ${error?.stack || error}`); process.exit(1); }); diff --git a/slab/bin/slab-repair-codex-watchers.mjs b/slab/bin/slab-repair-codex-watchers.mjs new file mode 100755 index 0000000000..1f004683d6 --- /dev/null +++ b/slab/bin/slab-repair-codex-watchers.mjs @@ -0,0 +1,32 @@ +#!/usr/bin/env node +// Reattach status watchers to live wrappers without restarting any Codex TUI. +import { readdirSync, readFileSync } from "node:fs"; +import { execFileSync, spawn } from "node:child_process"; +import { homedir } from "node:os"; +import { join, basename, dirname } from "node:path"; +import { fileURLToPath } from "node:url"; + +const root = process.env.SLAB_HOME || join(homedir(), ".local/share/slab"); +const active = join(root, "state/active-prompts"); +const watcher = join(dirname(fileURLToPath(import.meta.url)), "codex-session-watch.mjs"); +const processes = execFileSync("/bin/ps", ["-axo", "pid=,args="], { encoding: "utf8" }) + .split("\n").map(line => line.trim().split(/\s+/)); +let repaired = 0; +for (const sid of readdirSync(active)) { + try { + const marker = JSON.parse(readFileSync(join(active, sid), "utf8")); + if (marker.agent_type !== "codex" || !Number.isInteger(marker.wrapper_pid) || marker.wrapper_pid <= 0) continue; + process.kill(marker.wrapper_pid, 0); + for (const [pid, node, script, session] of processes) { + if (basename(node || "") === "node" && script?.endsWith("/codex-session-watch.mjs") && session === sid) { + try { process.kill(Number(pid), "SIGTERM"); } catch {} + } + } + spawn(process.execPath, [watcher, sid, String(Math.floor(Date.now()/1000)), + String(marker.wrapper_pid), marker.tty || "", marker.cwd || ""], { + detached: true, stdio: "ignore", env: process.env, + }).unref(); + repaired++; + } catch { /* Skip stale or partially written markers. */ } +} +console.log(`Reattached ${repaired} Codex status watchers`); diff --git a/slab/menubar-swift/Sources/SlabMenubar/ClaudeSession.swift b/slab/menubar-swift/Sources/SlabMenubar/ClaudeSession.swift index d62ea8eb9a..51db7f96fd 100644 --- a/slab/menubar-swift/Sources/SlabMenubar/ClaudeSession.swift +++ b/slab/menubar-swift/Sources/SlabMenubar/ClaudeSession.swift @@ -283,7 +283,7 @@ enum ClaudeSessionReader { // has overwritten it yet) — preserve so applyTerminalDecor // paints the appearance-matched bg. // (no-op: keep s.state == .blank) - } else if s.agentType == "easel", + } else if s.agentType == "easel" || s.agentType == "codex", s.state == .complete || s.state == .awaiting || s.state == .interrupted { // Easel receives app-server lifecycle events // directly, so its marker can state this transition without @@ -416,7 +416,7 @@ enum ClaudeSessionReader { case "blank": return .blank case "complete" where agentType == "easel": return .complete case "awaiting" where agentType == "easel": return .awaiting - case "interrupted" where agentType == "easel": return .interrupted + case "interrupted" where agentType == "easel" || agentType == "codex": return .interrupted default: return .working } }() diff --git a/slab/menubar-swift/deploy-host.sh b/slab/menubar-swift/deploy-host.sh index bab6da66e8..41446dd235 100755 --- a/slab/menubar-swift/deploy-host.sh +++ b/slab/menubar-swift/deploy-host.sh @@ -29,6 +29,7 @@ scp -q "$HERE/install.sh" "$HERE/Info.plist" \ "$HERE/computer.slab.menubar.plist.tmpl" "$TARGET:$STAGE/menubar-swift/" scp -q "$REPO/slab/bin/build-lock.sh" "$TARGET:$STAGE/bin/" scp -q "$REPO/slab/bin/codex-slab" "$REPO/slab/bin/slab-prepare-terminal.mjs" "$TARGET:$STAGE/bin/" +scp -q "$REPO/slab/bin/codex-session-watch.mjs" "$REPO/slab/bin/slab-repair-codex-watchers.mjs" "$TARGET:$STAGE/bin/" if [[ -f "$HERE/AppIcon.icns" ]]; then scp -q "$HERE/AppIcon.icns" "$TARGET:$STAGE/menubar-swift/" fi diff --git a/slab/menubar-swift/install.sh b/slab/menubar-swift/install.sh index c8f20a8160..e245ba3e27 100755 --- a/slab/menubar-swift/install.sh +++ b/slab/menubar-swift/install.sh @@ -217,12 +217,15 @@ provision_iterm2_profiles() { # Keep the launcher and its color-probe handoff with this installed release. provision_terminal_launcher() { mkdir -p "${HOME}/.local/bin" - for script in codex-slab slab-prepare-terminal.mjs; do + for script in codex-slab slab-prepare-terminal.mjs codex-session-watch.mjs slab-repair-codex-watchers.mjs; do local source="${SCRIPT_DIR}/../bin/${script}" [[ -f "$source" ]] || continue chmod +x "$source" ln -sfn "$source" "${HOME}/.local/bin/${script}" done + if command -v node >/dev/null 2>&1 && [[ -d "${HOME}/.local/share/slab/state/active-prompts" ]]; then + node "${SCRIPT_DIR}/../bin/slab-repair-codex-watchers.mjs" + fi } # Disable Terminal.app's "Do you want to terminate the running processes?" modal diff --git a/slab/test/codex-session-watch.test.mjs b/slab/test/codex-session-watch.test.mjs new file mode 100644 index 0000000000..6069bd410c --- /dev/null +++ b/slab/test/codex-session-watch.test.mjs @@ -0,0 +1,36 @@ +import test from "node:test"; +import assert from "node:assert/strict"; +import { handleLine, rolloutFromTitle } from "../bin/codex-session-watch.mjs"; + +const event = (type, ctx) => handleLine(JSON.stringify({ type: "event_msg", payload: { type } }), ctx); +test("long Codex turns keep the working heartbeat until completion, approval or interruption", () => { + for (const end of ["task_complete", "turn_complete", "turn_aborted", "approval_request"]) { + const ctx = { pending: [], lastUser: "test", turnActive: false }; + event("task_started", ctx); + assert.equal(ctx.turnActive, true); + handleLine(JSON.stringify({ type: "response_item", payload: { type: "function_call" } }), ctx); + assert.equal(ctx.turnActive, true, "Tool calls must not end the active turn"); + event(end, ctx); + assert.equal(ctx.turnActive, false, end); + event("task_started", ctx); + assert.equal(ctx.turnActive, true, "A subsequent turn returns to green"); + } +}); +test("reattaching publishes only the final status, without replaying historical work", () => { + const ctx = { pending: [], lastUser: "", turnActive: false, replaying: true }; + for (const type of ["task_started", "task_complete", "task_started"]) event(type, ctx); + assert.equal(ctx.pending.length, 1); + assert.equal(ctx.turnActive, true); +}); +test("shared daemon threads bind by this window's exact name and cwd, never recency", () => { + const rows = [ + { name: "Fix colors", cwd: "/work/repo", rollout_path: "/sessions/ours.jsonl" }, + { name: "Fix status", cwd: "/work/repo", rollout_path: "/sessions/peer.jsonl" }, + { name: "Fix colors", cwd: "/other/repo", rollout_path: "/sessions/other.jsonl" }, + ]; + assert.equal(rolloutFromTitle(rows, "◌ Fix colors | repo", "/work/repo"), rows[0].rollout_path); + assert.equal(rolloutFromTitle(rows, "◌ Fix status | repo", "/work/repo"), rows[1].rollout_path); + assert.equal(rolloutFromTitle(rows, "New session | repo", "/work/repo"), null); + assert.equal(rolloutFromTitle([...rows, {...rows[0], rollout_path: "/sessions/duplicate.jsonl"}], + "◌ Fix colors | repo", "/work/repo"), null); +}); -- 2.51.2