mirror of
https://github.com/OpenHands/OpenHands.git
synced 2026-10-07 16:19:05 +08:00
feat: execute git-info probe commands over bash-events WebSocket (#646)
Replace per-poll REST calls in useLocalGitInfo with a persistent WebSocket connection to /sockets/bash-events. Previously, useQuery's refetchInterval:10_000 caused AgentServerRuntimeService.executeCommand to fire individual HTTP requests on every tick (find + 2 git commands × up to 3 candidate dirs) to probe the workspace for git metadata. Changes: - websocket-url.ts: export buildBashWebSocketUrl() using the same host/pathPrefix extraction as the conversation-events URL builder - use-bash-command-runner.ts: new hook that maintains a persistent WS connection to /sockets/bash-events and exposes runCommand(command, cwd, timeout) → Promise. Commands queued during CONNECTING are flushed on open; all in-flight commands are rejected on close/error/unmount. - use-local-git-info.ts: swap AgentServerRuntimeService.executeCommand for useBashCommandRunner; keep refetchInterval:10_000 so branch changes (e.g. git checkout) are still reflected without a full refresh, but each probe now reuses the open socket rather than opening new HTTP connections. Co-authored-by: openhands <openhands@all-hands.dev>
This commit is contained in:
committed by
GitHub
co-authored by
openhands
parent
87b1079388
commit
2d494a4767
@@ -1,11 +1,11 @@
|
||||
import { useQuery } from "@tanstack/react-query";
|
||||
import { useRef } from "react";
|
||||
|
||||
import AgentServerRuntimeService, {
|
||||
CommandResult,
|
||||
} from "#/api/runtime-service/agent-server-runtime-service";
|
||||
import type { CommandResult } from "#/api/runtime-service/agent-server-runtime-service";
|
||||
import { getAgentServerWorkingDir } from "#/api/agent-server-config";
|
||||
import { useActiveConversation } from "#/hooks/query/use-active-conversation";
|
||||
import { useRuntimeIsReady } from "#/hooks/use-runtime-is-ready";
|
||||
import { useBashCommandRunner } from "#/hooks/use-bash-command-runner";
|
||||
import { Provider } from "#/types/settings";
|
||||
import { parseGitRemoteUrl } from "#/utils/parse-git-remote-url";
|
||||
|
||||
@@ -83,14 +83,19 @@ async function probeNestedRepoInDir(
|
||||
}
|
||||
|
||||
/**
|
||||
* Probe git metadata directly from the workspace checkout via the agent server
|
||||
* (`git remote get-url origin`, `git rev-parse --abbrev-ref HEAD`).
|
||||
* Probe git metadata directly from the workspace checkout via the agent
|
||||
* server's bash-events WebSocket (`git remote get-url origin`,
|
||||
* `git rev-parse --abbrev-ref HEAD`).
|
||||
*
|
||||
* We intentionally keep this probe enabled until the active conversation has
|
||||
* a complete repo tuple (`selected_repository`, `git_provider`,
|
||||
* `selected_branch`) so the control bar can recover from partial metadata
|
||||
* hydration after connect/clone flows.
|
||||
*
|
||||
* Commands are executed over a persistent WebSocket connection
|
||||
* (`/sockets/bash-events`) rather than individual REST calls, which avoids
|
||||
* the per-request HTTP overhead of the previous polling approach.
|
||||
*
|
||||
* Returns `null` fields when the working dir is not a git checkout — callers
|
||||
* should treat that the same as "no repo detected".
|
||||
*/
|
||||
@@ -107,6 +112,31 @@ export const useLocalGitInfo = () => {
|
||||
const hasConversationProvider = !!conversation?.git_provider;
|
||||
const hasConversationBranch = !!conversation?.selected_branch;
|
||||
|
||||
const queryEnabled =
|
||||
runtimeIsReady &&
|
||||
!!conversationId &&
|
||||
(!hasConversationRepo ||
|
||||
!hasConversationProvider ||
|
||||
!hasConversationBranch);
|
||||
|
||||
// Persistent WebSocket connection to the bash-events endpoint. The
|
||||
// connection is opened when the query is enabled and closed on unmount or
|
||||
// when the conversation changes.
|
||||
const runCommand = useBashCommandRunner(
|
||||
conversationUrl,
|
||||
sessionApiKey,
|
||||
queryEnabled,
|
||||
);
|
||||
|
||||
// Keep a ref so queryFn can call the latest runner without capturing it
|
||||
// as a queryKey dependency (runCommand is stable but the linter can't
|
||||
// infer that).
|
||||
const runCommandRef = useRef(runCommand);
|
||||
runCommandRef.current = runCommand;
|
||||
|
||||
// runCommandRef is a ref (always stable); the linter cannot infer this so
|
||||
// we disable the exhaustive-deps check here.
|
||||
// eslint-disable-next-line @tanstack/query/exhaustive-deps
|
||||
return useQuery<LocalGitInfo>({
|
||||
queryKey: [
|
||||
"local-git-info",
|
||||
@@ -117,13 +147,7 @@ export const useLocalGitInfo = () => {
|
||||
],
|
||||
queryFn: async () => {
|
||||
const run: RunCommand = (command, cwd, timeout) =>
|
||||
AgentServerRuntimeService.executeCommand(
|
||||
conversationUrl,
|
||||
sessionApiKey,
|
||||
command,
|
||||
cwd,
|
||||
timeout,
|
||||
);
|
||||
runCommandRef.current(command, cwd, timeout);
|
||||
const directInfo = await probeGitInfoAtDir(run, workingDir);
|
||||
if (directInfo.repository || directInfo.branch) return directInfo;
|
||||
|
||||
@@ -134,17 +158,13 @@ export const useLocalGitInfo = () => {
|
||||
|
||||
return EMPTY_LOCAL_GIT_INFO;
|
||||
},
|
||||
enabled:
|
||||
runtimeIsReady &&
|
||||
!!conversationId &&
|
||||
(!hasConversationRepo ||
|
||||
!hasConversationProvider ||
|
||||
!hasConversationBranch),
|
||||
enabled: queryEnabled,
|
||||
retry: false,
|
||||
// Re-probe the workspace every 10s so the UI reflects branch/repo
|
||||
// changes (e.g. `git checkout`, adding a remote) without requiring a
|
||||
// manual refresh when there is no `selected_repository` recorded on
|
||||
// the conversation.
|
||||
// the conversation. Commands now run over the persistent WebSocket
|
||||
// connection rather than individual REST calls.
|
||||
staleTime: 10_000,
|
||||
refetchInterval: 10_000,
|
||||
gcTime: 1000 * 60 * 5,
|
||||
|
||||
@@ -0,0 +1,191 @@
|
||||
import { useCallback, useEffect, useRef } from "react";
|
||||
import type {
|
||||
BashCommand,
|
||||
BashError,
|
||||
BashEvent,
|
||||
BashOutput,
|
||||
} from "@openhands/typescript-client";
|
||||
import type { CommandResult } from "#/api/runtime-service/agent-server-runtime-service";
|
||||
import { buildBashWebSocketUrl } from "#/utils/websocket-url";
|
||||
|
||||
interface WaitingCommand {
|
||||
command: string;
|
||||
cwd: string;
|
||||
timeout: number;
|
||||
resolve: (result: CommandResult) => void;
|
||||
reject: (error: Error) => void;
|
||||
}
|
||||
|
||||
interface PendingCommand {
|
||||
resolve: (result: CommandResult) => void;
|
||||
reject: (error: Error) => void;
|
||||
}
|
||||
|
||||
interface ActiveCommand extends PendingCommand {
|
||||
stdout: string[];
|
||||
stderr: string[];
|
||||
}
|
||||
|
||||
export type BashCommandRunner = (
|
||||
command: string,
|
||||
cwd: string,
|
||||
timeout: number,
|
||||
) => Promise<CommandResult>;
|
||||
|
||||
function isBashCommand(event: BashEvent): event is BashCommand {
|
||||
return event.kind === "BashCommand";
|
||||
}
|
||||
|
||||
function isBashOutput(event: BashEvent): event is BashOutput {
|
||||
return event.kind === "BashOutput";
|
||||
}
|
||||
|
||||
function isBashError(event: BashEvent): event is BashError {
|
||||
return event.kind === "BashError";
|
||||
}
|
||||
|
||||
/**
|
||||
* Maintains a persistent WebSocket connection to the agent-server's
|
||||
* `/sockets/bash-events` endpoint and exposes a `runCommand` function that
|
||||
* executes a bash command and returns a Promise that resolves when the
|
||||
* final `BashOutput` (non-null `exit_code`) arrives.
|
||||
*
|
||||
* Commands are correlated using a FIFO queue: each `BashCommand` echo
|
||||
* received from the server is paired with the oldest outstanding request in
|
||||
* the queue, and subsequent `BashOutput` events are matched by `command_id`.
|
||||
*
|
||||
* Commands queued while the socket is still in the CONNECTING state are
|
||||
* buffered and flushed automatically when the socket opens.
|
||||
*/
|
||||
export function useBashCommandRunner(
|
||||
conversationUrl: string | null | undefined,
|
||||
sessionApiKey: string | null | undefined,
|
||||
enabled: boolean,
|
||||
): BashCommandRunner {
|
||||
const wsRef = useRef<WebSocket | null>(null);
|
||||
// Commands waiting for the socket to transition from CONNECTING → OPEN
|
||||
const connectingQueueRef = useRef<WaitingCommand[]>([]);
|
||||
// Commands whose request was sent; waiting for the BashCommand echo to get command_id
|
||||
const pendingQueueRef = useRef<PendingCommand[]>([]);
|
||||
// Commands whose command_id is known; waiting for BashOutput with non-null exit_code
|
||||
const activeCommandsRef = useRef<Map<string, ActiveCommand>>(new Map());
|
||||
|
||||
useEffect(() => {
|
||||
if (!enabled) return;
|
||||
|
||||
const wsUrl = buildBashWebSocketUrl(conversationUrl, sessionApiKey);
|
||||
const ws = new WebSocket(wsUrl);
|
||||
wsRef.current = ws;
|
||||
|
||||
ws.onopen = () => {
|
||||
// Flush any commands that arrived while connecting
|
||||
for (const {
|
||||
command,
|
||||
cwd,
|
||||
timeout,
|
||||
resolve,
|
||||
reject,
|
||||
} of connectingQueueRef.current) {
|
||||
pendingQueueRef.current.push({ resolve, reject });
|
||||
ws.send(JSON.stringify({ command, cwd, timeout }));
|
||||
}
|
||||
connectingQueueRef.current = [];
|
||||
};
|
||||
|
||||
ws.onmessage = (event: MessageEvent) => {
|
||||
let data: BashEvent;
|
||||
try {
|
||||
data = JSON.parse(event.data as string) as BashEvent;
|
||||
} catch {
|
||||
return; // ignore malformed frames
|
||||
}
|
||||
|
||||
if (isBashCommand(data)) {
|
||||
// Associate the next pending request with the server-assigned command_id
|
||||
const pending = pendingQueueRef.current.shift();
|
||||
if (pending) {
|
||||
activeCommandsRef.current.set(data.id, {
|
||||
...pending,
|
||||
stdout: [],
|
||||
stderr: [],
|
||||
});
|
||||
}
|
||||
} else if (isBashOutput(data) && data.command_id) {
|
||||
const active = activeCommandsRef.current.get(data.command_id);
|
||||
if (active) {
|
||||
if (data.stdout) active.stdout.push(data.stdout);
|
||||
if (data.stderr) active.stderr.push(data.stderr);
|
||||
if (data.exit_code != null) {
|
||||
activeCommandsRef.current.delete(data.command_id);
|
||||
active.resolve({
|
||||
exit_code: data.exit_code,
|
||||
stdout: active.stdout.join(""),
|
||||
stderr: active.stderr.join(""),
|
||||
});
|
||||
}
|
||||
}
|
||||
} else if (isBashError(data)) {
|
||||
rejectAll(`Bash error: ${data.code}: ${data.detail}`);
|
||||
}
|
||||
};
|
||||
|
||||
function rejectAll(reason: string): void {
|
||||
const err = new Error(reason);
|
||||
for (const { reject: rej } of connectingQueueRef.current) rej(err);
|
||||
connectingQueueRef.current = [];
|
||||
for (const p of pendingQueueRef.current) p.reject(err);
|
||||
pendingQueueRef.current = [];
|
||||
for (const a of activeCommandsRef.current.values()) a.reject(err);
|
||||
activeCommandsRef.current.clear();
|
||||
}
|
||||
|
||||
ws.onclose = () => {
|
||||
wsRef.current = null;
|
||||
rejectAll("Bash WebSocket closed");
|
||||
};
|
||||
|
||||
ws.onerror = () => {
|
||||
wsRef.current = null;
|
||||
rejectAll("Bash WebSocket error");
|
||||
};
|
||||
|
||||
return () => {
|
||||
// Prevent the close/error handlers from double-rejecting after unmount
|
||||
ws.onclose = null;
|
||||
ws.onerror = null;
|
||||
ws.close();
|
||||
wsRef.current = null;
|
||||
rejectAll("Bash WebSocket unmounted");
|
||||
};
|
||||
}, [enabled, conversationUrl, sessionApiKey]);
|
||||
|
||||
const runCommand: BashCommandRunner = useCallback(
|
||||
(command: string, cwd: string, timeout: number) =>
|
||||
new Promise<CommandResult>((resolve, reject) => {
|
||||
const ws = wsRef.current;
|
||||
if (
|
||||
!ws ||
|
||||
ws.readyState === WebSocket.CLOSED ||
|
||||
ws.readyState === WebSocket.CLOSING
|
||||
) {
|
||||
reject(new Error("Bash WebSocket not available"));
|
||||
return;
|
||||
}
|
||||
if (ws.readyState === WebSocket.CONNECTING) {
|
||||
connectingQueueRef.current.push({
|
||||
command,
|
||||
cwd,
|
||||
timeout,
|
||||
resolve,
|
||||
reject,
|
||||
});
|
||||
} else {
|
||||
pendingQueueRef.current.push({ resolve, reject });
|
||||
ws.send(JSON.stringify({ command, cwd, timeout }));
|
||||
}
|
||||
}),
|
||||
[],
|
||||
);
|
||||
|
||||
return runCommand;
|
||||
}
|
||||
@@ -75,6 +75,34 @@ function getConversationUrlProtocol(
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Builds the WebSocket URL for the agent-server's bash-events endpoint.
|
||||
* The URL is derived from the same host and path prefix as the conversation
|
||||
* events socket so it works in both direct-connect and reverse-proxy deployments.
|
||||
*
|
||||
* @param conversationUrl The conversation URL containing host/port
|
||||
* @param sessionApiKey Optional session API key (appended as query param)
|
||||
* @returns WebSocket URL for the bash-events endpoint
|
||||
*/
|
||||
export function buildBashWebSocketUrl(
|
||||
conversationUrl: string | null | undefined,
|
||||
sessionApiKey?: string | null,
|
||||
): string {
|
||||
const baseHost = extractBaseHost(conversationUrl);
|
||||
const pathPrefix = extractPathPrefix(conversationUrl);
|
||||
|
||||
const pageIsSecure = window.location.protocol === "https:";
|
||||
const targetIsSecure =
|
||||
getConversationUrlProtocol(conversationUrl) === "https:";
|
||||
const protocol = pageIsSecure || targetIsSecure ? "wss:" : "ws:";
|
||||
|
||||
const base = `${protocol}//${baseHost}${pathPrefix}/sockets/bash-events`;
|
||||
if (sessionApiKey) {
|
||||
return `${base}?session_api_key=${encodeURIComponent(sessionApiKey)}`;
|
||||
}
|
||||
return base;
|
||||
}
|
||||
|
||||
/**
|
||||
* Builds the WebSocket URL for V1 conversations (without query params)
|
||||
* @param conversationId The conversation ID
|
||||
|
||||
Reference in New Issue
Block a user