open-swe/apps/cli/src/streaming.ts

238 lines
7.5 KiB
TypeScript
Raw Normal View History

feat: add updates to cli application (#503) * open swe cli: initial cli updates * openswe cli: add working authentication * open swe cli: remove unecessary console log statements * openswe cli: typescript modifications * openswe cli: updates to formats * openswe cli: updates to formatting * openswe cli: ran prettier * openswe cli: update urls * openswe cli: update eslint * openswe cli: add env example * openswe cli: add updates * openswe cli: formatting * openswe cli: add app installation * openswe cli: clean up lint and formatting * temp logs * l * working human feedback prompting * stop the re-renders * NOT WORKING -- KICKOFF PROGRAMMER * approve programmer * detected programmer working * some sembleance of logging * some changes * working logs * concatenated tool results * ui improvements * update ui * formatting * simple interrupt working * working multiple sessions + padding * working dissapearance of the approve/deny * TEST PUSH * cr * generally working tool calls * access thread ip properly * fix naming of env example * add readme * formatting * linting + formatting * dist updates * updates to yarn * no tests for cli * updates to readme * format * update readme * remove auth print out * add loading indicators * formatting changes * remove shared workspace * add lock * add dependency * update yarn lock * add back the installation token * replace with variable constants * Update apps/cli/src/auth-server.ts Co-authored-by: Brace Sproul <braceasproul@gmail.com> * openswe cli: updates * openswe cli: cleanup lint * openswe cli: update stream mode * openswe cli: update stream mode * openswe cli: streaming service * openswe cli: create selector for approval * openswe cli: remove streaming service * openswe cli: update types to langchain * openswe cli: update all streamModes * openswe cli: update formatting * openswe cli: update stream mode * openswe cli: notify of in development * openswe cli: change to istoolmessage --------- Co-authored-by: bracesproul <braceasproul@gmail.com>
2025-07-29 17:37:13 -04:00
import { Client, StreamMode } from "@langchain/langgraph-sdk";
import { v4 as uuidv4 } from "uuid";
import { encryptSecret } from "@open-swe/shared/crypto";
import {
MANAGER_GRAPH_ID,
GITHUB_TOKEN_COOKIE,
GITHUB_INSTALLATION_TOKEN_COOKIE,
GITHUB_INSTALLATION_NAME,
GITHUB_INSTALLATION_ID,
OPEN_SWE_STREAM_MODE,
} from "@open-swe/shared/constants";
import {
getAccessToken,
getInstallationAccessToken,
getInstallationId,
} from "./auth-server.js";
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 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
setClient: (client: Client) => void; // eslint-disable-line no-unused-vars
setThreadId: (id: string) => void; // eslint-disable-line no-unused-vars
}
export class StreamingService {
private callbacks: StreamingCallbacks;
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]);
}
}
}
}
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") {
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
}
}
}
this.callbacks.setStreamingPhase("done");
}
async startNewSession(prompt: string, selectedRepo: any) {
this.callbacks.setLogs(() => []);
this.callbacks.setLoadingLogs(true);
try {
const userAccessToken = getAccessToken();
const installationAccessToken = await getInstallationAccessToken();
const encryptionKey = process.env.SECRETS_ENCRYPTION_KEY;
if (!userAccessToken || !installationAccessToken || !encryptionKey) {
this.callbacks.setLogs(() => [
`Missing secrets: ${userAccessToken ? "" : "userAccessToken, "}${installationAccessToken ? "" : "installationAccessToken, "}${encryptionKey ? "" : "encryptionKey"}`,
]);
return;
}
const encryptedUserToken = encryptSecret(userAccessToken, encryptionKey);
const encryptedInstallationToken = encryptSecret(
installationAccessToken,
encryptionKey,
);
const [owner, repoName] = selectedRepo.full_name.split("/");
const runInput = {
messages: [
{
id: uuidv4(),
type: "human",
content: [{ type: "text", text: prompt }],
},
],
targetRepository: {
owner,
repo: repoName,
branch: selectedRepo.default_branch || "main",
},
autoAcceptPlan: false,
};
const installationId = getInstallationId();
const newClient = new Client({
apiUrl: LANGGRAPH_URL,
defaultHeaders: {
[GITHUB_TOKEN_COOKIE]: encryptedUserToken,
[GITHUB_INSTALLATION_TOKEN_COOKIE]: encryptedInstallationToken,
[GITHUB_INSTALLATION_NAME]: owner,
[GITHUB_INSTALLATION_ID]: installationId,
},
});
this.callbacks.setClient(newClient);
const thread = await newClient.threads.create();
const threadId = thread.thread_id;
this.callbacks.setThreadId(threadId);
const run = await newClient.runs.create(threadId, MANAGER_GRAPH_ID, {
input: runInput,
config: { recursion_limit: 400 },
ifNotExists: "create",
streamResumable: true,
streamMode: OPEN_SWE_STREAM_MODE as StreamMode[],
});
await this.startManagerStream(newClient, threadId, run.run_id);
} catch (err: any) {
this.callbacks.setLogs((prev) => [
...prev,
`Error during streaming: ${err.message}`,
]);
this.callbacks.setLoadingLogs(false);
} finally {
this.callbacks.setLoadingLogs(false);
}
}
}