mirror of
https://github.com/Sea-Haven-Industries/open-swe.git
synced 2026-09-30 08:03:15 +00:00
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
This commit is contained in:
parent
53e5dd954d
commit
3caa492ffd
7 changed files with 718 additions and 581 deletions
3
.gitignore
vendored
3
.gitignore
vendored
|
|
@ -47,3 +47,6 @@ credentials.json
|
|||
.langgraph_api
|
||||
|
||||
**/.claude/settings.local.json
|
||||
|
||||
# Test traces
|
||||
apps/cli/test_traces/
|
||||
|
|
|
|||
|
|
@ -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 <file>", "Replay from LangSmith trace file")
|
||||
.option("--speed <ms>", "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 (
|
||||
<Box justifyContent="center" paddingY={2}>
|
||||
<Text>
|
||||
{text}
|
||||
{dots}
|
||||
</Text>
|
||||
</Box>
|
||||
);
|
||||
};
|
||||
// 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<string[]>([]);
|
||||
const [plannerFeedback, setPlannerFeedback] = useState<string | null>(null);
|
||||
const [streamingPhase, setStreamingPhase] = useState<
|
||||
"streaming" | "awaitingFeedback" | "done"
|
||||
>("streaming");
|
||||
const [plannerThreadId, setPlannerThreadId] = useState<string | null>(null);
|
||||
const [hasStartedChat, setHasStartedChat] = useState(false);
|
||||
const [loadingLogs, setLoadingLogs] = useState(false);
|
||||
const [logs, setLogs] = useState<string[]>([]);
|
||||
const [streamingService, setStreamingService] =
|
||||
useState<StreamingService | null>(null);
|
||||
const [currentInterrupt, setCurrentInterrupt] = useState<{
|
||||
command: string;
|
||||
args: Record<string, any>;
|
||||
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 (
|
||||
<Box flexDirection="row" alignItems="center" gap={2}>
|
||||
<Text>Plan feedback: </Text>
|
||||
<Text
|
||||
color={selectedOption === "approve" ? "black" : "white"}
|
||||
bold={selectedOption === "approve"}
|
||||
>
|
||||
{selectedOption === "approve" ? "▶ " : " "}Approve
|
||||
</Text>
|
||||
<Text
|
||||
color={selectedOption === "deny" ? "black" : "white"}
|
||||
bold={selectedOption === "deny"}
|
||||
>
|
||||
{selectedOption === "deny" ? "▶ " : " "}Deny
|
||||
</Text>
|
||||
<Text dimColor>(Use ←/→ to select, Enter to confirm)</Text>
|
||||
</Box>
|
||||
);
|
||||
};
|
||||
|
||||
// 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 (
|
||||
<Box flexDirection="column" height={process.stdout.rows}>
|
||||
{/* Auto-scrolling logs area - strict boundary container */}
|
||||
<Box
|
||||
height={availableLogHeight}
|
||||
flexDirection="column"
|
||||
paddingX={1}
|
||||
paddingBottom={1}
|
||||
overflowY="hidden"
|
||||
flexShrink={0}
|
||||
justifyContent="flex-end"
|
||||
>
|
||||
<Box flexDirection="column">
|
||||
{loadingLogs && logs.length === 0 ? (
|
||||
<LoadingSpinner text="Starting agent" />
|
||||
) : (
|
||||
visibleLogs.map((log, index) => (
|
||||
<Box key={`${logs.length}-${index}`}>
|
||||
<Text
|
||||
dimColor={
|
||||
!log.startsWith("[AI]") && !log.includes("PROPOSED PLAN")
|
||||
}
|
||||
bold={log.startsWith("[AI]") || log.includes("PROPOSED PLAN")}
|
||||
>
|
||||
{log}
|
||||
</Text>
|
||||
</Box>
|
||||
))
|
||||
)}
|
||||
</Box>
|
||||
</Box>
|
||||
|
||||
{/* Welcome message right above input bar */}
|
||||
{/* Welcome message or logs display */}
|
||||
{!hasStartedChat ? (
|
||||
<Box flexDirection="column" paddingX={1}>
|
||||
<Box>
|
||||
|
|
@ -237,13 +129,91 @@ const App: React.FC = () => {
|
|||
## ## ## ## ## ## ## #### ## ######### ## ## ## ## ## ##
|
||||
## ######### ## #### ## ## ## ## ## ######### ## ## ####
|
||||
## ## ## ## ### ## ## ## ## ## ## ## ## ## ## ###
|
||||
######## ## ## ## ## ###### ###### ## ## ## ## #### ## ## OPEN SWE CLI
|
||||
######## ## ## ## ## ###### ###### ## ## ## ## #### ## ##
|
||||
`}
|
||||
</Text>
|
||||
</Box>
|
||||
</Box>
|
||||
) : (
|
||||
<Box height={8} />
|
||||
<Box
|
||||
flexDirection="column"
|
||||
height={availableHeight}
|
||||
paddingX={2}
|
||||
paddingY={1}
|
||||
paddingBottom={3}
|
||||
>
|
||||
<Box
|
||||
flexDirection="column"
|
||||
height={availableHeight - 5}
|
||||
justifyContent="flex-end"
|
||||
overflow="hidden"
|
||||
>
|
||||
{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 (
|
||||
<Box
|
||||
key={index}
|
||||
paddingLeft={isToolCall ? 1 : isToolResult ? 2 : 0}
|
||||
width="100%"
|
||||
flexShrink={0}
|
||||
>
|
||||
<Text
|
||||
color={
|
||||
isAIMessage
|
||||
? "magenta"
|
||||
: isToolResult
|
||||
? "gray"
|
||||
: isRemovedLine
|
||||
? "redBright"
|
||||
: isAddedLine
|
||||
? "greenBright"
|
||||
: isLongBashCommand
|
||||
? "gray"
|
||||
: undefined
|
||||
}
|
||||
bold={isAIMessage}
|
||||
wrap="wrap"
|
||||
>
|
||||
{log}
|
||||
</Text>
|
||||
</Box>
|
||||
);
|
||||
})}
|
||||
</Box>
|
||||
</Box>
|
||||
)}
|
||||
|
||||
{/* Approval prompt above input when interrupt is active */}
|
||||
{currentInterrupt && (
|
||||
<Box paddingX={2} paddingY={1}>
|
||||
<Text color="magenta">
|
||||
Approve this command? $ {currentInterrupt.command}{" "}
|
||||
{currentInterrupt.args.path ||
|
||||
Object.values(currentInterrupt.args).join(" ")}{" "}
|
||||
(yes/no/custom)
|
||||
</Text>
|
||||
</Box>
|
||||
)}
|
||||
|
||||
{/* Cooking icon above input when loading */}
|
||||
{loadingLogs && (
|
||||
<Box paddingX={2} paddingY={1}>
|
||||
<Text>Thinking...</Text>
|
||||
</Box>
|
||||
)}
|
||||
|
||||
{/* Fixed input area at bottom */}
|
||||
|
|
@ -257,28 +227,39 @@ const App: React.FC = () => {
|
|||
justifyContent="center"
|
||||
>
|
||||
<Box>
|
||||
{streamingPhase === "awaitingFeedback" ? (
|
||||
<PlannerFeedbackInput />
|
||||
) : !hasStartedChat ? (
|
||||
{replayFile ? (
|
||||
<Text>> Replay mode - input disabled</Text>
|
||||
) : (
|
||||
<CustomInput
|
||||
onSubmit={(value) => {
|
||||
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);
|
||||
}
|
||||
}}
|
||||
/>
|
||||
) : (
|
||||
<Box>
|
||||
<Text>Streaming...</Text>
|
||||
</Box>
|
||||
)}
|
||||
</Box>
|
||||
</Box>
|
||||
|
|
@ -286,7 +267,7 @@ const App: React.FC = () => {
|
|||
{/* Local mode indicator underneath the input bar */}
|
||||
<Box paddingX={2} paddingY={0}>
|
||||
<Text>
|
||||
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
|
||||
</Text>
|
||||
</Box>
|
||||
</Box>
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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<string, string | number | boolean>;
|
||||
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);
|
||||
|
|
|
|||
100
apps/cli/src/trace_replay.ts
Normal file
100
apps/cli/src/trace_replay.ts
Normal file
|
|
@ -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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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";
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue