From 3caa492ffd92ee50fa97b8452c09b5993cd9bdf5 Mon Sep 17 00:00:00 2001 From: Palash Shah <35114859+Palashio@users.noreply.github.com> Date: Fri, 22 Aug 2025 16:58:27 -0400 Subject: [PATCH] feat: modify cli for deep agents + formatting (#792) * feat: update cli for local * feat: file diff formatting * fix: udate diff format * fix: show green/red diffs * fix: modifications * feat: str based replacement * feat: update streaming to allow submissions * feat: traces * fix: wrap text * feat: update cli * fix: remove json * fix: remove write fill diff * fix: mods * fix: updates to logging * fix: add trace replay service * feat: update logs * fix: remove debugging * fix: lint * fix: unused vars * fix: updates * fix: formatting * fix: update streaming value * fix: replace constants * feat: mods * nit: change interrupt * nit: remove question --- .gitignore | 3 + apps/cli/src/index.tsx | 313 ++++++++++----------- apps/cli/src/logger.ts | 461 ++++++++++++++++--------------- apps/cli/src/streaming.ts | 371 ++++++++++++++----------- apps/cli/src/trace_replay.ts | 100 +++++++ apps/cli/src/utils.ts | 50 ++-- packages/shared/src/constants.ts | 1 + 7 files changed, 718 insertions(+), 581 deletions(-) create mode 100644 apps/cli/src/trace_replay.ts diff --git a/.gitignore b/.gitignore index 69b3d6dc..7c3e716a 100644 --- a/.gitignore +++ b/.gitignore @@ -47,3 +47,6 @@ credentials.json .langgraph_api **/.claude/settings.local.json + +# Test traces +apps/cli/test_traces/ diff --git a/apps/cli/src/index.tsx b/apps/cli/src/index.tsx index 45de1b86..198d9918 100644 --- a/apps/cli/src/index.tsx +++ b/apps/cli/src/index.tsx @@ -3,22 +3,28 @@ import React, { useState, useEffect } from "react"; import { render, Box, Text, useInput } from "ink"; import { Command } from "commander"; import { OPEN_SWE_CLI_VERSION } from "./constants.js"; +import fs from "fs"; import dotenv from "dotenv"; dotenv.config(); +// Keep the process alive - prevents exit when streaming completes +const keepAlive = setInterval(() => {}, 60000); + // Handle graceful exit on Ctrl+C and Ctrl+K process.on("SIGINT", () => { + clearInterval(keepAlive); console.log("\nšŸ‘‹ Goodbye!"); process.exit(0); }); process.on("SIGTERM", () => { + clearInterval(keepAlive); console.log("\nšŸ‘‹ Goodbye!"); process.exit(0); }); -import { submitFeedback } from "./utils.js"; import { StreamingService } from "./streaming.js"; +import { TraceReplayService } from "./trace_replay.js"; // Parse command line arguments with Commander const program = new Command(); @@ -27,46 +33,21 @@ program .name("open-swe") .description("Open SWE CLI - Local Mode") .version(OPEN_SWE_CLI_VERSION) + .option("--replay ", "Replay from LangSmith trace file") + .option("--speed ", "Replay speed in milliseconds", "500") .helpOption("-h, --help", "Display help for command") .parse(); // Always run in local mode process.env.OPEN_SWE_LOCAL_MODE = "true"; -console.log("šŸ  Starting Open SWE CLI in Local Mode"); -console.log(" Working directory:", process.cwd()); -console.log(" No GitHub authentication required"); -console.log(""); - -const LoadingSpinner: React.FC<{ text: string }> = ({ text }) => { - const [dots, setDots] = useState(""); - - useEffect(() => { - const interval = setInterval(() => { - setDots((prev) => (prev.length >= 3 ? "" : prev + ".")); - }, 500); - return () => clearInterval(interval); - }, []); - - return ( - - - {text} - {dots} - - - ); -}; // eslint-disable-next-line no-unused-vars const CustomInput: React.FC<{ onSubmit: (value: string) => void }> = ({ onSubmit, }) => { const [input, setInput] = useState(""); - const [isSubmitted, setIsSubmitted] = useState(false); useInput((inputChar: string, key: { [key: string]: any }) => { - if (isSubmitted) return; - // Handle Ctrl+K for exit if (key.ctrl && inputChar.toLowerCase() === "k") { console.log("\nšŸ‘‹ Goodbye!"); @@ -76,13 +57,9 @@ const CustomInput: React.FC<{ onSubmit: (value: string) => void }> = ({ if (key.return) { if (input.trim()) { // Only submit if there's actual content - setIsSubmitted(true); onSubmit(input); - // Reset for next input - setTimeout(() => { - setInput(""); - setIsSubmitted(false); - }, 100); + // Clear input immediately after submission + setInput(""); } } else if (key.backspace || key.delete) { setInput((prev) => prev.slice(0, -1)); @@ -99,132 +76,47 @@ const CustomInput: React.FC<{ onSubmit: (value: string) => void }> = ({ }; const App: React.FC = () => { - const [logs, setLogs] = useState([]); - const [plannerFeedback, setPlannerFeedback] = useState(null); - const [streamingPhase, setStreamingPhase] = useState< - "streaming" | "awaitingFeedback" | "done" - >("streaming"); - const [plannerThreadId, setPlannerThreadId] = useState(null); const [hasStartedChat, setHasStartedChat] = useState(false); const [loadingLogs, setLoadingLogs] = useState(false); + const [logs, setLogs] = useState([]); + const [streamingService, setStreamingService] = + useState(null); + const [currentInterrupt, setCurrentInterrupt] = useState<{ + command: string; + args: Record; + id: string; + } | null>(null); - const PlannerFeedbackInput: React.FC = () => { - const [selectedOption, setSelectedOption] = useState< - "approve" | "deny" | null - >(null); + const options = program.opts(); + const replayFile = options.replay; + const playbackSpeed = parseInt(options.speed) || 500; - useInput((inputChar: string, key: { [key: string]: any }) => { - if (streamingPhase !== "awaitingFeedback") return; - - // Handle Ctrl+K for exit - if (key.ctrl && inputChar.toLowerCase() === "k") { - console.log("\nšŸ‘‹ Goodbye!"); - process.exit(0); - } - - if (key.return && selectedOption) { - setPlannerFeedback(selectedOption); - setSelectedOption(null); - } else if (key.leftArrow) { - setSelectedOption("approve"); - } else if (key.rightArrow) { - setSelectedOption("deny"); - } - }); - - if (streamingPhase !== "awaitingFeedback") { - return null; - } - - return ( - - Plan feedback: - - {selectedOption === "approve" ? "ā–¶ " : " "}Approve - - - {selectedOption === "deny" ? "ā–¶ " : " "}Deny - - (Use ←/→ to select, Enter to confirm) - - ); - }; - - // Handle planner feedback + // Auto-start replay if file provided useEffect(() => { - if ( - streamingPhase === "awaitingFeedback" && - plannerFeedback && - plannerThreadId - ) { - (async () => { - await submitFeedback({ - plannerFeedback, - plannerThreadId, + if (replayFile && !hasStartedChat) { + try { + const traceData = JSON.parse(fs.readFileSync(replayFile, "utf8")); + setHasStartedChat(true); + + const traceReplayService = new TraceReplayService({ setLogs, - setPlannerFeedback: () => setPlannerFeedback(null), - setStreamingPhase, + setLoadingLogs, }); - })(); + + traceReplayService.replayFromTrace(traceData, playbackSpeed); + } catch (err: any) { + console.error("Error loading replay file:", err.message); + process.exit(1); + } } - }, [streamingPhase, plannerFeedback, plannerThreadId]); + }, [replayFile, hasStartedChat, playbackSpeed]); - const headerHeight = 0; const inputHeight = 4; - const welcomeHeight = hasStartedChat ? 0 : 8; - const paddingHeight = 3; - const availableLogHeight = Math.max( - 5, - process.stdout.rows - - headerHeight - - inputHeight - - welcomeHeight - - paddingHeight, - ); - - // Always show the most recent logs (auto-scroll to bottom) - const visibleLogs = - logs.length > availableLogHeight ? logs.slice(-availableLogHeight) : logs; + const availableHeight = process.stdout.rows - inputHeight - 1; return ( - {/* Auto-scrolling logs area - strict boundary container */} - - - {loadingLogs && logs.length === 0 ? ( - - ) : ( - visibleLogs.map((log, index) => ( - - - {log} - - - )) - )} - - - - {/* Welcome message right above input bar */} + {/* Welcome message or logs display */} {!hasStartedChat ? ( @@ -237,13 +129,91 @@ const App: React.FC = () => { ## ## ## ## ## ## ## #### ## ######### ## ## ## ## ## ## ## ######### ## #### ## ## ## ## ## ######### ## ## #### ## ## ## ## ### ## ## ## ## ## ## ## ## ## ## ### -######## ## ## ## ## ###### ###### ## ## ## ## #### ## ## OPEN SWE CLI +######## ## ## ## ## ###### ###### ## ## ## ## #### ## ## `} ) : ( - + + + {logs + .filter( + (log) => + log !== null && log !== undefined && typeof log === "string", + ) + .map((log, index) => { + const isToolCall = log.startsWith("ā–ø"); + const isToolResult = log.startsWith(" ↳"); + const isAIMessage = log.startsWith("ā—†"); + const isRemovedLine = log.startsWith("- "); + const isAddedLine = log.startsWith("+ "); + const isLongBashCommand = + isToolCall && + (log.includes("execute_bash:") || log.includes("shell:")) && + log.includes("..."); + + return ( + + + {log} + + + ); + })} + + + )} + + {/* Approval prompt above input when interrupt is active */} + {currentInterrupt && ( + + + Approve this command? $ {currentInterrupt.command}{" "} + {currentInterrupt.args.path || + Object.values(currentInterrupt.args).join(" ")}{" "} + (yes/no/custom) + + + )} + + {/* Cooking icon above input when loading */} + {loadingLogs && ( + + Thinking... + )} {/* Fixed input area at bottom */} @@ -257,28 +227,39 @@ const App: React.FC = () => { justifyContent="center" > - {streamingPhase === "awaitingFeedback" ? ( - - ) : !hasStartedChat ? ( + {replayFile ? ( + > Replay mode - input disabled + ) : ( { - setHasStartedChat(true); - setPlannerFeedback(null); + // Handle interrupt approval responses + if (currentInterrupt && streamingService) { + streamingService.submitInterruptResponse(value); + return; + } - const streamingService = new StreamingService({ - setLogs, - setPlannerThreadId, - setStreamingPhase, - setLoadingLogs, - }); + if (!streamingService) { + // First message - create new session + setHasStartedChat(true); + // Clear logs only for first message + setLogs([]); - streamingService.startNewSession(value); + const newStreamingService = new StreamingService({ + setLogs, + setLoadingLogs, + setCurrentInterrupt, + setStreamingPhase: () => {}, + }); + + setStreamingService(newStreamingService); + newStreamingService.startNewSession(value); + } else { + // If stream is active, submit to existing stream + // If stream is not active, also submit to existing stream + streamingService.submitToExistingStream(value); + } }} /> - ) : ( - - Streaming... - )} @@ -286,7 +267,7 @@ const App: React.FC = () => { {/* Local mode indicator underneath the input bar */} - Working on {process.env.OPEN_SWE_LOCAL_PROJECT_PATH} • Ctrl+K to exit + Working on {process.env.OPEN_SWE_LOCAL_PROJECT_PATH} • Ctrl+C to exit diff --git a/apps/cli/src/logger.ts b/apps/cli/src/logger.ts index 47511337..266ab976 100644 --- a/apps/cli/src/logger.ts +++ b/apps/cli/src/logger.ts @@ -22,6 +22,43 @@ interface LogChunk { ops?: Array<{ value: string }>; } +/** + * Create a simple diff between old and new strings + */ +function createSimpleDiff(oldString: string, newString: string): string[] { + const logs: string[] = []; + + if (!oldString && newString) { + const lines = newString.split("\n").slice(0, 10); + lines.forEach((line) => logs.push(`+ ${line}`)); + if (newString.split("\n").length > 10) { + logs.push(`+ ... (${newString.split("\n").length - 10} more lines)`); + } + return logs; + } + + if (!newString) { + oldString.split("\n").forEach((line) => logs.push(`- ${line}`)); + return logs; + } + + const oldLines = oldString.split("\n"); + const newLines = newString.split("\n"); + + const removedLines = oldLines.filter( + (oldLine) => !newLines.some((newLine) => newLine === oldLine), + ); + + const addedLines = newLines.filter( + (newLine) => !oldLines.some((oldLine) => oldLine === newLine), + ); + + removedLines.forEach((line) => logs.push(`- ${line}`)); + addedLines.forEach((line) => logs.push(`+ ${line}`)); + + return logs; +} + /** * Format a tool call arguments into a clean, readable string */ @@ -31,22 +68,71 @@ function formatToolCallArgs(tool: ToolCall): string { if (!tool.args) return toolName; switch (toolName.toLowerCase()) { - case "shell": { + case "shell": + case "execute_bash": { + let command = ""; if (Array.isArray(tool.args.command)) { - return `${toolName}: ${tool.args.command.join(" ")}`; + command = tool.args.command.join(" "); + } else { + command = tool.args.command || ""; } - return `${toolName}: ${tool.args.command || ""}`; + + // Truncate long commands (more than 160 characters) + if (command.length > 160) { + return `${toolName}: ${command.substring(0, 160)}...`; + } + return `${toolName}: ${command}`; } - case "grep": { + case "write_file": { + const filePath = tool.args.file_path || ""; + const content = tool.args.content || ""; + const lineCount = content.split("\n").length; + return `${toolName}: ${filePath} (${lineCount} lines)`; + } + + case "read_file": { + const filePath = tool.args.file_path || ""; + return `${toolName}: ${filePath}`; + } + + case "edit_file": { + const filePath = tool.args.file_path || ""; + return `${toolName}: ${filePath}`; + } + + case "http_request": { + const method = tool.args.method || "GET"; + const url = tool.args.url || ""; + return `${toolName}: ${method} ${url}`; + } + + case "web_search": { const query = tool.args.query || ""; return `${toolName}: "${query}"`; } + case "grep": { + const pattern = tool.args.pattern || ""; + const path = tool.args.path || ""; + return `${toolName}: "${pattern}"${path ? ` in ${path}` : ""}`; + } + + case "glob": { + const pattern = tool.args.pattern || ""; + const path = tool.args.path || ""; + return `${toolName}: ${pattern}${path ? ` in ${path}` : ""}`; + } + case "view": { return `${toolName}: ${tool.args.path || ""}`; } + case "ls": { + const path = tool.args.path || ""; + return `${toolName}: ${path}`; + } + case "str_replace_based_edit_tool": { const command = tool.args.command || ""; @@ -57,9 +143,7 @@ function formatToolCallArgs(tool: ToolCall): string { return `${toolName}: insert_line=${insertLine}, new_str="${newStr}"`; } case "str_replace": { - const oldStr = tool.args.old_str || ""; - const newStr = tool.args.new_str || ""; - return `${toolName}: old_str="${oldStr}", new_str="${newStr}"`; + return `${toolName}: string replacement`; } case "create": { const fileText = tool.args.file_text || ""; @@ -77,118 +161,24 @@ function formatToolCallArgs(tool: ToolCall): string { } } - case "search_documents_for": { - const query = tool.args.query || ""; - const url = tool.args.url || ""; - return `${toolName}: "${query}" in ${url}`; - } - - case "get_url_content": { - return `${toolName}: ${tool.args.url || ""}`; - } - - case "session_plan": { - const title = tool.args.title || ""; - const planSteps = tool.args.plan || []; - if (title) { - return `${toolName}: "${title}" (${planSteps.length} steps)`; + case "write_todos": { + const todos = tool.args.todos || []; + if (Array.isArray(todos)) { + const todoCount = todos.length; + const statusCounts = todos.reduce((acc: any, todo: any) => { + acc[todo.status] = (acc[todo.status] || 0) + 1; + return acc; + }, {}); + const statusSummary = Object.entries(statusCounts) + .map(([status, count]) => `${count} ${status}`) + .join(", "); + return `${toolName}: Updated ${todoCount} todos (${statusSummary})`; } - return `${toolName}: ${planSteps.length} plan steps`; - } - - case "apply_patch": { - const filePath = tool.args.file_path || ""; - const diff = tool.args.diff || ""; - const diffLines = diff.split("\n").length; - return `${toolName}: applied ${diffLines} line diff to ${filePath}`; - } - - case "install_dependencies": { - const command = tool.args.command || []; - if (Array.isArray(command)) { - return `${toolName}: ${command.join(" ")}`; - } - return `${toolName}: ${command}`; - } - - case "scratchpad": { - const scratchpad = tool.args.scratchpad || []; - if (Array.isArray(scratchpad)) { - return `${toolName}: ${scratchpad.length} notes`; - } - return `${toolName}: ${scratchpad}`; - } - - case "command_safety_evaluator": { - const command = tool.args.command || ""; - return `${toolName}: evaluating "${command}"`; - } - - case "respond_and_route": { - const response = tool.args.response || ""; - const route = tool.args.route || ""; - if (response && route) { - return `${toolName}: "${response}" → ${route}`; - } else if (response) { - return `${toolName}: "${response}"`; - } else if (route) { - return `${toolName}: → ${route}`; - } - return `${toolName}: routing decision`; - } - - case "request_human_help": { - const helpRequest = tool.args.help_request || ""; - return `${toolName}: "${helpRequest}"`; - } - - case "update_plan": { - const reasoning = tool.args.update_plan_reasoning || ""; - return `${toolName}: ${reasoning.slice(0, 50)}...`; - } - - case "mark_task_completed": { - const summary = tool.args.completed_task_summary || ""; - return `${toolName}: ${summary.slice(0, 50)}...`; - } - - case "mark_task_not_completed": { - const reasoning = tool.args.reasoning || ""; - return `${toolName}: ${reasoning.slice(0, 50)}...`; - } - - case "diagnose_error": { - const diagnosis = tool.args.diagnosis || ""; - return `${toolName}: ${diagnosis.slice(0, 50)}...`; - } - - case "write_technical_notes": { - const notes = tool.args.notes || ""; - return `${toolName}: ${notes.slice(0, 50)}...`; - } - - case "summarize_conversation_history": { - const reasoning = tool.args.reasoning || ""; - return `${toolName}: ${reasoning.slice(0, 50)}...`; - } - - case "code_review_mark_task_completed": { - const review = tool.args.review || ""; - return `${toolName}: ${review.slice(0, 50)}...`; - } - - case "code_review_mark_task_not_complete": { - const review = tool.args.review || ""; - const actions = tool.args.additional_actions || []; - return `${toolName}: ${review.slice(0, 30)}... (${actions.length} actions)`; - } - - case "review_started": { - const started = tool.args.review_started || false; - return `${toolName}: ${started ? "started" : "not started"}`; + return `${toolName}: Updated todos`; } } - return ""; + + return toolName; } /** @@ -207,7 +197,55 @@ function formatToolResult(message: ToolMessage): string { switch (toolName.toLowerCase()) { case "shell": - return content; + case "execute_bash": { + try { + const result = JSON.parse(content); + if (!result.success && result.stderr) { + return result.stderr; + } + if (result.success && result.stdout) { + return result.stdout; + } + return content; + } catch { + return content; + } + } + + case "write_file": + if (isError) return content; + + return "File written successfully"; + + case "read_file": { + const contentLength = content.length; + return `${contentLength} characters`; + } + + case "edit_file": + return isError ? content : "File edited successfully"; + + case "http_request": { + try { + const result = JSON.parse(content); + return `HTTP ${result.status_code || "unknown"}: ${result.success ? "Success" : "Failed"}`; + } catch { + return content.length > 100 ? content.slice(0, 100) + "..." : content; + } + } + + case "web_search": { + try { + const result = JSON.parse(content); + if (result.error) { + return `Search error: ${result.error}`; + } + const results = result.results || []; + return `${results.length} search results found`; + } catch { + return content.length > 100 ? content.slice(0, 100) + "..." : content; + } + } case "grep": { if (content.includes("Exit code 1. No results found.")) { @@ -219,9 +257,7 @@ function formatToolResult(message: ToolMessage): string { case "view": { const contentLength = content.length; - return contentLength > 1000 - ? `${contentLength} characters (truncated)` - : `${contentLength} characters`; + return `${contentLength} characters`; } case "str_replace_based_edit_tool": @@ -230,19 +266,22 @@ function formatToolResult(message: ToolMessage): string { case "get_url_content": return `${content.length} characters of content`; - case "apply_patch": - return "Patch applied successfully"; - - case "install_dependencies": - return "Dependencies installed successfully"; - - case "command_safety_evaluator": - try { - const evaluation = JSON.parse(content); - return `Safety: ${evaluation.is_safe ? "SAFE" : "UNSAFE"} (${evaluation.risk_level} risk)`; - } catch { - return content; + case "write_todos": + if (content.includes("Updated todo list")) { + return "Todo list updated successfully"; } + return content.length > 100 ? content.slice(0, 100) + "..." : content; + + case "ls": + try { + const items = JSON.parse(content); + if (Array.isArray(items)) { + return `${items.length} items: ${items.slice(0, 8).join(", ")}${items.length > 8 ? "..." : ""}`; + } + } catch { + // fallthrough to default + } + return content.length > 100 ? content.slice(0, 100) + "..." : content; default: return content.length > 200 ? content.slice(0, 200) + "..." : content; @@ -251,52 +290,12 @@ function formatToolResult(message: ToolMessage): string { export function formatDisplayLog(chunk: LogChunk | string): string[] { if (typeof chunk === "string") { - if (chunk.startsWith("Human feedback:")) { - return [ - `[HUMAN FEEDBACK RECEIVED] ${chunk.replace("Human feedback:", "").trim()}`, - ]; - } - if (chunk.startsWith("Interrupt:")) { - const message = chunk.replace("Interrupt:", "").trim(); - return [ - "═══════════════════════════════════════", - `šŸ“¤ INTERRUPT: "${message}"`, - "═══════════════════════════════════════", - ]; - } - // Filter out raw file content and object references - if ( - chunk === "[object Object]" || - chunk.includes("total 4") || - chunk.includes("drwxr-xr-x") || - chunk.includes("Exit code 1") || - chunk.startsWith("#") || - chunk.startsWith("-") || - chunk.startsWith("./") - ) { - return []; - } - // Single line system messages - const cleanChunk = chunk.replace(/\s+/g, " ").trim(); - const maxLength = 150; - const truncated = - cleanChunk.length > maxLength - ? cleanChunk.slice(0, maxLength) + "... [trunc]" - : cleanChunk; - return [`[SYSTEM] ${truncated}`]; + return [chunk]; } const data = chunk.data; const logs: string[] = []; - // Handle session events - if (data.plannerSession) { - logs.push("[PLANNER SESSION STARTED]"); - } - if (data.programmerSession) { - logs.push("[PROGRAMMER SESSION STARTED]"); - } - // Handle messages const nestedDataObj = Object.values(data)[0] as unknown as Record< string, @@ -317,16 +316,17 @@ export function formatDisplayLog(chunk: LogChunk | string): string[] { // Handle tool messages if (isToolMessage(message)) { const toolName = message.name || "tool"; + + // Skip displaying results for todo list tool calls + if (toolName === "write_todos") { + continue; + } + const result = formatToolResult(message); if (result) { - // Concatenate long tool results to a single line (truncate if too long) - const maxLength = 500; + // Display tool results as indented subsections let formattedResult = result.replace(/\s+/g, " "); - if (formattedResult.length > maxLength) { - formattedResult = - formattedResult.slice(0, maxLength) + "... [trunc]"; - } - logs.push(`[TOOL RESULT] ${toolName}: ${formattedResult}`); + logs.push(` ↳ ${formattedResult}`); } continue; } @@ -338,12 +338,7 @@ export function formatDisplayLog(chunk: LogChunk | string): string[] { const reasoning = String(message.additional_kwargs.reasoning) .replace(/\s+/g, " ") .trim(); - const maxLength = 150; - const truncated = - reasoning.length > maxLength - ? reasoning.slice(0, maxLength) + "... [trunc]" - : reasoning; - logs.push(`[REASONING] ${truncated}`); + logs.push(`[REASONING] ${reasoning}`); } // Handle tool calls @@ -353,7 +348,50 @@ export function formatDisplayLog(chunk: LogChunk | string): string[] { message.tool_calls.forEach((tool) => { const formattedArgs = formatToolCallArgs(tool); - logs.push(`[TOOL CALL] ${formattedArgs}`); + logs.push(`ā–ø ${formattedArgs}`); + + // Special handling for write_todos to display the actual todos nicely + if ( + tool.name === "write_todos" && + tool.args && + tool.args.todos && + Array.isArray(tool.args.todos) + ) { + const todos = tool.args.todos; + logs.push(""); // blank line before todos + todos.forEach((todo: any) => { + const statusIcon = + todo.status === "completed" + ? "āœ“" + : todo.status === "in_progress" + ? "→" + : "ā—‹"; + logs.push(` ${statusIcon} ${todo.content}`); + }); + } + + // Special handling for edit_file to display the diff + if (tool.name === "edit_file" && tool.args) { + const oldString = tool.args.old_string || ""; + const newString = tool.args.new_string || ""; + const diffLines = createSimpleDiff(oldString, newString); + logs.push(...diffLines); + } + + // Special handling for write_file to display the new content + if (tool.name === "write_file" && tool.args) { + const content = tool.args.content || ""; + const diffLines = createSimpleDiff("", content); + logs.push(...diffLines); + } + + // Special handling for str_replace_based_edit_tool to display the diff + if (tool.name === "str_replace_based_edit_tool" && tool.args) { + const oldStr = tool.args.old_str || ""; + const newStr = tool.args.new_str || ""; + const diffLines = createSimpleDiff(oldStr, newStr); + logs.push(...diffLines); + } // Handle technical notes from tool call if ( @@ -376,14 +414,9 @@ export function formatDisplayLog(chunk: LogChunk | string): string[] { // Handle regular AI messages const text = getMessageContentString(message.content); if (text) { - // Always single line, remove newlines and truncate + // Always single line, remove newlines const cleanText = text.replace(/\s+/g, " ").trim(); - const maxLength = 200; - const truncated = - cleanText.length > maxLength - ? cleanText.slice(0, maxLength) + "... [trunc]" - : cleanText; - logs.push(`[AI] ${truncated}`); + logs.push(`ā—† ${cleanText}`); } } @@ -393,12 +426,7 @@ export function formatDisplayLog(chunk: LogChunk | string): string[] { if (text) { // Single line human messages const cleanText = text.replace(/\s+/g, " ").trim(); - const maxLength = 150; - const truncated = - cleanText.length > maxLength - ? cleanText.slice(0, maxLength) + "... [trunc]" - : cleanText; - logs.push(`[HUMAN] ${truncated}`); + logs.push(`ā—‰ ${cleanText}`); } } } catch (error: any) { @@ -406,31 +434,27 @@ export function formatDisplayLog(chunk: LogChunk | string): string[] { // Fallback to original message if conversion fails if (msg.type === "tool") { const toolName = msg.name || "tool"; - const content = getMessageContentString(msg.content); - if (content) { - logs.push(`[TOOL RESULT] ${toolName}: ${content}`); + + // Skip displaying results for todo list tool calls + if (toolName === "write_todos") { + // Skip this tool result + } else { + const content = getMessageContentString(msg.content); + if (content) { + logs.push(` ↳ ${content}`); + } } } else if (msg.type === "ai") { const text = getMessageContentString(msg.content); if (text) { const cleanText = text.replace(/\s+/g, " ").trim(); - const maxLength = 200; - const truncated = - cleanText.length > maxLength - ? cleanText.slice(0, maxLength) + "... [trunc]" - : cleanText; - logs.push(`[AI] ${truncated}`); + logs.push(`ā—† ${cleanText}`); } } else if (msg.type === "human") { const text = getMessageContentString(msg.content); if (text) { const cleanText = text.replace(/\s+/g, " ").trim(); - const maxLength = 150; - const truncated = - cleanText.length > maxLength - ? cleanText.slice(0, maxLength) + "... [trunc]" - : cleanText; - logs.push(`[HUMAN] ${truncated}`); + logs.push(`ā—‰ ${cleanText}`); } } } @@ -459,16 +483,11 @@ export function formatDisplayLog(chunk: LogChunk | string): string[] { "šŸŽÆ PROPOSED PLAN", ...steps.map((step: string, idx: number) => ` ${idx + 1}. ${step}`), - " ", // Blank line after - ); - } else { - logs.push( - " ", // Blank line for separation - "ā³ INTERRUPT: Waiting for feedback...", " ", // Blank line after ); } } + return logs; } diff --git a/apps/cli/src/streaming.ts b/apps/cli/src/streaming.ts index 51069981..76c61e58 100644 --- a/apps/cli/src/streaming.ts +++ b/apps/cli/src/streaming.ts @@ -1,204 +1,255 @@ import { Client, StreamMode } from "@langchain/langgraph-sdk"; -import { v4 as uuidv4 } from "uuid"; -import { - MANAGER_GRAPH_ID, - LOCAL_MODE_HEADER, - OPEN_SWE_STREAM_MODE, -} from "@open-swe/shared/constants"; +import { LOCAL_MODE_HEADER } from "@open-swe/shared/constants"; import { formatDisplayLog } from "./logger.js"; -import { isAgentInboxInterruptSchema } from "@open-swe/shared/agent-inbox-interrupt"; const LANGGRAPH_URL = process.env.LANGGRAPH_URL || "http://localhost:2024"; +interface InterruptData { + command: string; + args: Record; + id: string; +} + +interface InterruptItem { + id: string; + value: InterruptData; +} + +interface StreamChunk { + event: string; + data: ChunkData; +} + +interface ChunkData { + __interrupt__?: InterruptItem[]; + agent?: { + messages: Array<{ + role: string; + content: string; + }>; + }; + [key: string]: unknown; +} + interface StreamingCallbacks { setLogs: (updater: (prev: string[]) => string[]) => void; // eslint-disable-line no-unused-vars - setPlannerThreadId: (id: string) => void; // eslint-disable-line no-unused-vars - setStreamingPhase: (phase: "streaming" | "awaitingFeedback" | "done") => void; // eslint-disable-line no-unused-vars setLoadingLogs: (loading: boolean) => void; // eslint-disable-line no-unused-vars + setCurrentInterrupt: (interrupt: InterruptData | null) => void; // eslint-disable-line no-unused-vars + setStreamingPhase: (phase: string) => void; // eslint-disable-line no-unused-vars } export class StreamingService { private callbacks: StreamingCallbacks; + private client: Client | null = null; + private threadId: string | null = null; + private rawLogs: (string | StreamChunk)[] = []; constructor(callbacks: StreamingCallbacks) { this.callbacks = callbacks; } - private async handleProgrammerStream( - client: Client, - programmerThreadId: string, - programmerRunId: string, - ) { - for await (const programmerChunk of client.runs.joinStream( - programmerThreadId, - programmerRunId, - { - streamMode: OPEN_SWE_STREAM_MODE as StreamMode[], - }, - )) { - if (programmerChunk.event === "updates") { - const formatted = formatDisplayLog(programmerChunk); - if (formatted.length > 0) { - this.callbacks.setLogs((prev) => [...prev, ...formatted]); - } - } - } - } + /** + * Get formatted logs for display + */ + getFormattedLogs(): string[] { + const formattedLogs: string[] = []; - private async handlePlannerStream( - client: Client, - plannerThreadId: string, - plannerRunId: string, - ): Promise<{ needsFeedback: boolean }> { - let programmerStreamed = false; - - for await (const subChunk of client.runs.joinStream( - plannerThreadId, - plannerRunId, - { - streamMode: OPEN_SWE_STREAM_MODE as StreamMode[], - }, - )) { - if (subChunk.event === "updates") { - const formatted = formatDisplayLog(subChunk); - // Filter out human messages from planner stream (already logged in manager) - const filteredFormatted = formatted.filter( - (log) => !log.startsWith("[HUMAN]"), - ); - if (filteredFormatted.length > 0) { - this.callbacks.setLogs((prev) => [...prev, ...filteredFormatted]); - } - } - - // Check for programmer session - if ( - !programmerStreamed && - subChunk.data?.programmerSession?.threadId && - typeof subChunk.data.programmerSession.threadId === "string" && - typeof subChunk.data.programmerSession.runId === "string" - ) { - programmerStreamed = true; - await this.handleProgrammerStream( - client, - subChunk.data.programmerSession.threadId, - subChunk.data.programmerSession.runId, - ); - } - - // Detect HumanInterrupt in planner stream - const interruptArr = - subChunk.data && Array.isArray(subChunk.data["__interrupt__"]) - ? subChunk.data["__interrupt__"] - : undefined; - const firstInterruptValue = - interruptArr && interruptArr[0] && interruptArr[0].value - ? interruptArr[0].value - : undefined; - - if (isAgentInboxInterruptSchema(firstInterruptValue)) { - return { needsFeedback: true }; - } - } - - return { needsFeedback: false }; - } - - private async startManagerStream( - client: Client, - threadId: string, - runId: string, - ) { - let plannerStreamed = false; - - for await (const chunk of client.runs.joinStream(threadId, runId)) { - if (chunk.event === "updates") { + for (const chunk of this.rawLogs) { + if (typeof chunk === "string") { const formatted = formatDisplayLog(chunk); - if (formatted.length > 0) { - this.callbacks.setLogs((prev) => { - if (prev.length === 0) { - this.callbacks.setLoadingLogs(false); - } - return [...prev, ...formatted]; - }); - } - } - - // Check for plannerSession - if ( - !plannerStreamed && - chunk.data && - chunk.data.plannerSession && - typeof chunk.data.plannerSession.threadId === "string" && - typeof chunk.data.plannerSession.runId === "string" - ) { - plannerStreamed = true; - this.callbacks.setPlannerThreadId(chunk.data.plannerSession.threadId); - - const result = await this.handlePlannerStream( - client, - chunk.data.plannerSession.threadId, - chunk.data.plannerSession.runId, - ); - - if (result.needsFeedback) { - this.callbacks.setStreamingPhase("awaitingFeedback"); - return; // Pause streaming, let React render feedback prompt - } + formattedLogs.push(...formatted); + } else if (chunk && chunk.data) { + // Process all chunks with data, not just "updates" events + const formatted = formatDisplayLog(chunk); + formattedLogs.push(...formatted); } } - this.callbacks.setStreamingPhase("done"); + return formattedLogs; } + /** + * Update the display with formatted logs + */ + private updateDisplay() { + const formattedLogs = this.getFormattedLogs(); + this.callbacks.setLogs(() => formattedLogs); + } + + /** + * Start a new session + */ async startNewSession(prompt: string) { + this.rawLogs = []; this.callbacks.setLogs(() => []); this.callbacks.setLoadingLogs(true); + // Keeping for the future, not needed now try { - const runInput = { - messages: [ - { - id: uuidv4(), - type: "human", - content: [{ type: "text", text: prompt }], - }, - ], - targetRepository: { - owner: "local", - repo: "local", - branch: "main", - }, - autoAcceptPlan: false, - }; - const headers = { [LOCAL_MODE_HEADER]: "true", }; - const newClient = new Client({ + this.client = new Client({ apiUrl: LANGGRAPH_URL, defaultHeaders: headers, }); - const thread = await newClient.threads.create(); - const threadId = thread.thread_id; + const thread = await this.client.threads.create(); + this.threadId = thread.thread_id; - const run = await newClient.runs.create(threadId, MANAGER_GRAPH_ID, { - input: runInput, - config: { - recursion_limit: 400, + // Stream using the pattern from deep-agents + const stream = await this.client.runs.stream(this.threadId, "coding", { + input: { + messages: [ + { + role: "system", + content: + "You are working on " + + (process.env.OPEN_SWE_LOCAL_PROJECT_PATH || ""), + }, + { role: "user", content: prompt }, + ], }, - ifNotExists: "create", - streamResumable: true, - streamMode: OPEN_SWE_STREAM_MODE as StreamMode[], + streamMode: ["updates"] as StreamMode[], }); - await this.startManagerStream(newClient, threadId, run.run_id); - } catch (err: any) { - this.callbacks.setLogs((prev) => [ - ...prev, - `Error during streaming: ${err.message}`, - ]); + // Process the stream + for await (const chunk of stream) { + this.updateDisplay(); + + if (chunk.event === "updates") { + // Check for interrupts in the chunk + if (chunk.data && chunk.data.__interrupt__) { + const chunkData = chunk.data as ChunkData; + const interrupt = chunkData.__interrupt__?.[0]?.value; + if (interrupt?.command && interrupt?.args) { + this.callbacks.setCurrentInterrupt({ + command: interrupt.command, + args: interrupt.args, + id: chunkData.__interrupt__?.[0]?.id || "unknown", + }); + } + } + + // Store raw chunk instead of formatting immediately + this.rawLogs.push(chunk); + this.updateDisplay(); + + if (this.rawLogs.length === 1) { + this.callbacks.setLoadingLogs(false); + } + } + } + + this.callbacks.setStreamingPhase("done"); + } catch (err: unknown) { + const errorMessage = err instanceof Error ? err.message : "Unknown error"; + this.rawLogs.push(`Error during streaming: ${errorMessage}`); + this.updateDisplay(); + this.callbacks.setLoadingLogs(false); + } finally { + this.callbacks.setLoadingLogs(false); + } + } + + async submitInterruptResponse(response: boolean | string) { + if (!this.client || !this.threadId) { + throw new Error("No active stream session. Start a new session first."); + } + + // Clear the interrupt from UI + this.callbacks.setCurrentInterrupt(null); + this.callbacks.setLoadingLogs(true); + + try { + const stream = await this.client.runs.stream(this.threadId, "coding", { + command: { resume: response }, + streamMode: ["updates"] as StreamMode[], + }); + + // Process the stream + for await (const chunk of stream) { + if (chunk.event === "updates") { + // Check for interrupts in the chunk + if (chunk.data && chunk.data.__interrupt__) { + const chunkData = chunk.data as ChunkData; + const interrupt = chunkData.__interrupt__?.[0]?.value; + if (interrupt?.command && interrupt?.args) { + this.callbacks.setCurrentInterrupt({ + command: interrupt.command, + args: interrupt.args, + id: chunkData.__interrupt__?.[0]?.id || "unknown", + }); + } + } + + // Store raw chunk instead of formatting immediately + this.rawLogs.push(chunk); + this.updateDisplay(); + + if (this.rawLogs.length === 1) { + this.callbacks.setLoadingLogs(false); + } + } + } + + this.callbacks.setStreamingPhase("done"); + } catch (err: unknown) { + const errorMessage = err instanceof Error ? err.message : "Unknown error"; + this.rawLogs.push(`Error submitting approval: ${errorMessage}`); + this.updateDisplay(); + this.callbacks.setLoadingLogs(false); + } finally { + this.callbacks.setLoadingLogs(false); + } + } + + async submitToExistingStream(prompt: string) { + if (!this.client || !this.threadId) { + throw new Error("No active stream session. Start a new session first."); + } + + // Don't clear logs - continue the conversation + this.callbacks.setLoadingLogs(true); + + try { + const stream = await this.client.runs.stream(this.threadId, "coding", { + input: { + messages: [{ role: "user", content: prompt }], + }, + streamMode: ["updates"] as StreamMode[], + }); + + // Process the stream + for await (const chunk of stream) { + if (chunk.event === "updates") { + // Check for interrupts in the chunk + if (chunk.data && chunk.data.__interrupt__) { + const chunkData = chunk.data as ChunkData; + const interrupt = chunkData.__interrupt__?.[0]?.value; + if (interrupt?.command && interrupt?.args) { + this.callbacks.setCurrentInterrupt({ + command: interrupt.command, + args: interrupt.args, + id: chunkData.__interrupt__?.[0]?.id || "unknown", + }); + } + } + + // Store raw chunk instead of formatting immediately + this.rawLogs.push(chunk); + this.updateDisplay(); + + if (this.rawLogs.length === 1) { + this.callbacks.setLoadingLogs(false); + } + } + } + } catch (err: unknown) { + const errorMessage = err instanceof Error ? err.message : "Unknown error"; + this.rawLogs.push(`Error submitting to stream: ${errorMessage}`); + this.updateDisplay(); this.callbacks.setLoadingLogs(false); } finally { this.callbacks.setLoadingLogs(false); diff --git a/apps/cli/src/trace_replay.ts b/apps/cli/src/trace_replay.ts new file mode 100644 index 00000000..151d172d --- /dev/null +++ b/apps/cli/src/trace_replay.ts @@ -0,0 +1,100 @@ +import { formatDisplayLog } from "./logger.js"; + +export interface TraceReplayCallbacks { + setLogs: (updater: (prev: string[]) => string[]) => void; // eslint-disable-line no-unused-vars + setLoadingLogs: (loading: boolean) => void; // eslint-disable-line no-unused-vars +} + +export class TraceReplayService { + private callbacks: TraceReplayCallbacks; + private rawLogs: any[] = []; + + constructor(callbacks: TraceReplayCallbacks) { + this.callbacks = callbacks; + } + + /** + * Get formatted logs for display + */ + getFormattedLogs(): string[] { + const formattedLogs: string[] = []; + + for (const chunk of this.rawLogs) { + if (typeof chunk === "string") { + const formatted = formatDisplayLog(chunk); + formattedLogs.push(...formatted); + } else if (chunk && chunk.data) { + // Process all chunks with data, not just "updates" events + const formatted = formatDisplayLog(chunk); + formattedLogs.push(...formatted); + } + } + + return formattedLogs; + } + + /** + * Update the display with formatted logs + */ + private updateDisplay() { + const formattedLogs = this.getFormattedLogs(); + this.callbacks.setLogs(() => formattedLogs); + } + + async replayFromTrace(langsmithRun: any, playbackSpeed: number = 500) { + this.rawLogs = []; + this.callbacks.setLogs(() => []); + this.callbacks.setLoadingLogs(true); + + try { + const messages = langsmithRun.messages || []; + + for (let i = 0; i < messages.length; i++) { + const message = messages[i]; + + // Convert LangSmith message to the format expected by formatDisplayLog + const mockChunk = { + event: "updates", + data: { + agent: { + messages: [message], + }, + }, + }; + this.rawLogs.push(mockChunk); + this.updateDisplay(); + + if (this.rawLogs.length === 1) { + this.callbacks.setLoadingLogs(false); + } + + // Add delay between messages to simulate streaming + if (i < messages.length - 1) { + await new Promise((resolve) => setTimeout(resolve, playbackSpeed)); + } + } + + // Check for interrupt data in the trace and add it at the end + if (langsmithRun.__interrupt__ || langsmithRun.interrupt) { + const interruptData = + langsmithRun.__interrupt__ || langsmithRun.interrupt; + const interruptChunk = { + event: "interrupt", + data: { + __interrupt__: Array.isArray(interruptData) + ? interruptData + : [interruptData], + }, + }; + this.rawLogs.push(interruptChunk); + this.updateDisplay(); + } + } catch (err: any) { + this.rawLogs.push(`Error during replay: ${err.message}`); + this.updateDisplay(); + this.callbacks.setLoadingLogs(false); + } finally { + this.callbacks.setLoadingLogs(false); + } + } +} diff --git a/apps/cli/src/utils.ts b/apps/cli/src/utils.ts index 515ef2f8..0c0c68c9 100644 --- a/apps/cli/src/utils.ts +++ b/apps/cli/src/utils.ts @@ -5,15 +5,15 @@ import { Client, StreamMode } from "@langchain/langgraph-sdk"; import { OPEN_SWE_STREAM_MODE, - PLANNER_GRAPH_ID, LOCAL_MODE_HEADER, + OPEN_SWE_V2_GRAPH_ID, } from "@open-swe/shared/constants"; import { formatDisplayLog } from "./logger.js"; const LANGGRAPH_URL = process.env.LANGGRAPH_URL || "http://localhost:2024"; /** - * Submit feedback to the planner + * Submit feedback to the coding agent */ export async function submitFeedback({ plannerFeedback, @@ -46,46 +46,28 @@ export async function submitFeedback({ } // Create a new stream with the feedback - const stream = await client.runs.stream(plannerThreadId, PLANNER_GRAPH_ID, { - command: { - resume: [ - { - type: plannerFeedback === "approve" ? "accept" : "ignore", - args: null, - }, - ], + const stream = await client.runs.stream( + plannerThreadId, + OPEN_SWE_V2_GRAPH_ID, + { + command: { + resume: [ + { + type: plannerFeedback === "approve" ? "accept" : "ignore", + args: null, + }, + ], + }, + streamMode: OPEN_SWE_STREAM_MODE as StreamMode[], }, - streamMode: OPEN_SWE_STREAM_MODE as StreamMode[], - }); + ); - let programmerStreamed = false; // Process the stream response for await (const chunk of stream) { const formatted = formatDisplayLog(chunk); if (formatted.length > 0) { setLogs((prev) => [...prev, ...formatted]); } - - // Check for programmer session in the resumed planner stream - const chunkData = chunk.data as any; - if ( - !programmerStreamed && - chunkData?.programmerSession?.threadId && - typeof chunkData.programmerSession.threadId === "string" && - typeof chunkData.programmerSession.runId === "string" - ) { - programmerStreamed = true; - // Join programmer stream - for await (const programmerChunk of client.runs.joinStream( - chunkData.programmerSession.threadId, - chunkData.programmerSession.runId, - )) { - const formatted = formatDisplayLog(programmerChunk); - if (formatted.length > 0) { - setLogs((prev) => [...prev, ...formatted]); - } - } - } } // Set streaming phase to done when complete diff --git a/packages/shared/src/constants.ts b/packages/shared/src/constants.ts index 36ee650c..16423ca2 100644 --- a/packages/shared/src/constants.ts +++ b/packages/shared/src/constants.ts @@ -17,6 +17,7 @@ export const GITHUB_AUTH_STATE_COOKIE = "github_auth_state"; export const GITHUB_INSTALLATION_ID_COOKIE = "github_installation_id"; export const GITHUB_TOKEN_TYPE_COOKIE = "github_token_type"; +export const OPEN_SWE_V2_GRAPH_ID = "open-swe-v2"; export const MANAGER_GRAPH_ID = "manager"; export const PLANNER_GRAPH_ID = "planner"; export const PROGRAMMER_GRAPH_ID = "programmer";