fix: RPC parity, UX polish, logging

This commit is contained in:
Cristina Poncela Cubeiro
2026-06-18 15:59:32 +02:00
parent c4e89b0337
commit 337de9b078
3 changed files with 56 additions and 20 deletions
+14 -6
View File
@@ -18,7 +18,7 @@ const packageJson = JSON.parse(readFileSync(join(__dirname, "../package.json"),
function printHelp(): void { function printHelp(): void {
console.log( console.log(
`orchestrator v${packageJson.version}\n\nUsage:\n orchestrator serve\n orchestrator list\n orchestrator spawn [--cwd <path>] [--label <label>]\n orchestrator status <instance-id>\n orchestrator stop <instance-id>\n orchestrator rpc <instance-id> <json-command>\n orchestrator attach <instance-id>\n orchestrator --help\n orchestrator --version`, `orchestrator v${packageJson.version}\n\nUsage:\n orchestrator serve\n orchestrator list\n orchestrator spawn [--cwd <path>] [--label <label>]\n orchestrator status <instance-id>\n orchestrator stop <instance-id>\n orchestrator rpc <instance-id> <json-command>\n orchestrator attach <instance-id>\n orchestrator --help\n orchestrator --version\n\nAttach stdin expects JSONL RpcCommand or extension_ui_response messages.`,
); );
} }
@@ -36,6 +36,7 @@ function getFlagValue(args: string[], flag: string): string | undefined {
async function attach(instanceId: string): Promise<void> { async function attach(instanceId: string): Promise<void> {
const socket = createConnection(getSocketPath()); const socket = createConnection(getSocketPath());
let stdinBuffer = "";
process.stdin.setEncoding("utf8"); process.stdin.setEncoding("utf8");
await new Promise<void>((resolve, reject) => { await new Promise<void>((resolve, reject) => {
@@ -49,6 +50,7 @@ async function attach(instanceId: string): Promise<void> {
socket.on("data", (chunk: Buffer | string) => { socket.on("data", (chunk: Buffer | string) => {
process.stdout.write(chunk.toString()); process.stdout.write(chunk.toString());
}); });
console.error(`attached to ${instanceId}; send JSONL RpcCommand or extension_ui_response on stdin`);
socket.on("error", (error) => { socket.on("error", (error) => {
console.error(error instanceof Error ? error.message : String(error)); console.error(error instanceof Error ? error.message : String(error));
process.exit(1); process.exit(1);
@@ -57,11 +59,17 @@ async function attach(instanceId: string): Promise<void> {
process.exit(0); process.exit(0);
}); });
process.stdin.on("data", (chunk: string) => { process.stdin.on("data", (chunk: string) => {
const lines = chunk stdinBuffer += chunk;
.split("\n") while (true) {
.map((line) => line.trim()) const newlineIndex = stdinBuffer.indexOf("\n");
.filter((line) => line.length > 0); if (newlineIndex === -1) {
for (const line of lines) { return;
}
const line = stdinBuffer.slice(0, newlineIndex).trim();
stdinBuffer = stdinBuffer.slice(newlineIndex + 1);
if (!line) {
continue;
}
const parsed = JSON.parse(line) as RpcCommand | RpcExtensionUIResponse; const parsed = JSON.parse(line) as RpcCommand | RpcExtensionUIResponse;
if (parsed.type === "extension_ui_response") { if (parsed.type === "extension_ui_response") {
socket.write(encodeMessage(parsed)); socket.write(encodeMessage(parsed));
+34 -8
View File
@@ -88,6 +88,22 @@ function computeBackoffDelayMs(failureCount: number): number {
return Math.min(HEARTBEAT_BACKOFF_MAX_MS, exponentialDelay + jitterMs); return Math.min(HEARTBEAT_BACKOFF_MAX_MS, exponentialDelay + jitterMs);
} }
function formatRadiusError(error: unknown): string {
if (error instanceof RadiusHttpError) {
return `HTTP ${error.status}: ${error.message}`;
}
if (error instanceof Error) {
return error.message;
}
return String(error);
}
function logRadiusRetry(scope: string, action: string, delayMs: number, failureCount: number, error: unknown): void {
console.error(
`${scope} ${action} failed (attempt ${failureCount}); retrying in ${delayMs}ms: ${formatRadiusError(error)}`,
);
}
export function getRadiusUrl(): string { export function getRadiusUrl(): string {
return process.env.PI_RADIUS_URL || DEFAULT_RADIUS_URL; return process.env.PI_RADIUS_URL || DEFAULT_RADIUS_URL;
} }
@@ -295,7 +311,7 @@ export class RadiusPresence {
if (!isNotFoundError(error)) { if (!isNotFoundError(error)) {
this.machineTransientFailureCount += 1; this.machineTransientFailureCount += 1;
const delayMs = computeBackoffDelayMs(this.machineTransientFailureCount); const delayMs = computeBackoffDelayMs(this.machineTransientFailureCount);
console.error(`Radius machine heartbeat failed; retrying in ${delayMs}ms`, error); logRadiusRetry("Radius machine", "heartbeat", delayMs, this.machineTransientFailureCount, error);
this.scheduleMachineHeartbeat(delayMs); this.scheduleMachineHeartbeat(delayMs);
return; return;
} }
@@ -312,7 +328,13 @@ export class RadiusPresence {
} catch (recoveryError) { } catch (recoveryError) {
this.machineTransientFailureCount += 1; this.machineTransientFailureCount += 1;
const delayMs = computeBackoffDelayMs(this.machineTransientFailureCount); const delayMs = computeBackoffDelayMs(this.machineTransientFailureCount);
console.error(`Radius machine re-registration failed; retrying in ${delayMs}ms`, recoveryError); logRadiusRetry(
"Radius machine",
"re-registration",
delayMs,
this.machineTransientFailureCount,
recoveryError,
);
this.scheduleMachineHeartbeat(delayMs); this.scheduleMachineHeartbeat(delayMs);
} }
} }
@@ -337,7 +359,7 @@ export class RadiusPresence {
if (!isNotFoundError(error)) { if (!isNotFoundError(error)) {
state.transientFailureCount += 1; state.transientFailureCount += 1;
const delayMs = computeBackoffDelayMs(state.transientFailureCount); const delayMs = computeBackoffDelayMs(state.transientFailureCount);
console.error(`Radius Pi heartbeat failed for instance ${instanceId}; retrying in ${delayMs}ms`, error); logRadiusRetry(`Radius Pi ${instanceId}`, "heartbeat", delayMs, state.transientFailureCount, error);
this.schedulePiHeartbeat(instanceId, delayMs); this.schedulePiHeartbeat(instanceId, delayMs);
return; return;
} }
@@ -352,14 +374,18 @@ export class RadiusPresence {
try { try {
const recovered = await this.reRegisterPi(instanceId); const recovered = await this.reRegisterPi(instanceId);
if (!recovered) { if (!recovered) {
console.error(`Radius Pi re-registration skipped for instance ${instanceId}`); const delayMs = computeBackoffDelayMs(1);
this.schedulePiHeartbeat(instanceId, computeBackoffDelayMs(1)); console.error(`Radius Pi ${instanceId} re-registration skipped; retrying in ${delayMs}ms`);
this.schedulePiHeartbeat(instanceId, delayMs);
} }
} catch (recoveryError) { } catch (recoveryError) {
state.transientFailureCount += 1; state.transientFailureCount += 1;
const delayMs = computeBackoffDelayMs(state.transientFailureCount); const delayMs = computeBackoffDelayMs(state.transientFailureCount);
console.error( logRadiusRetry(
`Radius Pi re-registration failed for instance ${instanceId}; retrying in ${delayMs}ms`, `Radius Pi ${instanceId}`,
"re-registration",
delayMs,
state.transientFailureCount,
recoveryError, recoveryError,
); );
this.schedulePiHeartbeat(instanceId, delayMs); this.schedulePiHeartbeat(instanceId, delayMs);
@@ -376,7 +402,7 @@ export class RadiusPresence {
try { try {
await this.reRegisterPi(instance.id); await this.reRegisterPi(instance.id);
} catch (error) { } catch (error) {
console.error(`Radius Pi re-registration failed for instance ${instance.id}`, error); console.error(`Radius Pi ${instance.id} re-registration failed: ${formatRadiusError(error)}`);
} }
} }
} }
+8 -6
View File
@@ -11,6 +11,10 @@ function error(id: string | undefined, command: string, message: string): RpcRes
return { id, type: "response", command, success: false, error: message }; return { id, type: "response", command, success: false, error: message };
} }
function assertNeverRpcCommand(command: never): never {
throw new Error(`Unknown command: ${(command as { type: string }).type}`);
}
export async function handleRpcCommand(runtime: AgentSessionRuntime, command: RpcCommand): Promise<RpcResponse> { export async function handleRpcCommand(runtime: AgentSessionRuntime, command: RpcCommand): Promise<RpcResponse> {
const session = runtime.session; const session = runtime.session;
const id = command.id; const id = command.id;
@@ -68,7 +72,7 @@ export async function handleRpcCommand(runtime: AgentSessionRuntime, command: Rp
case "fork": { case "fork": {
const result = await runtime.fork(command.entryId); const result = await runtime.fork(command.entryId);
return success(id, "fork", { text: result.selectedText ?? "", cancelled: result.cancelled }); return success(id, "fork", { text: result.selectedText, cancelled: result.cancelled });
} }
case "clone": { case "clone": {
@@ -186,7 +190,7 @@ export async function handleRpcCommand(runtime: AgentSessionRuntime, command: Rp
} }
case "get_last_assistant_text": { case "get_last_assistant_text": {
const text = session.getLastAssistantText() ?? null; const text = session.getLastAssistantText();
return success(id, "get_last_assistant_text", { text }); return success(id, "get_last_assistant_text", { text });
} }
@@ -236,9 +240,7 @@ export async function handleRpcCommand(runtime: AgentSessionRuntime, command: Rp
return success(id, "get_commands", { commands }); return success(id, "get_commands", { commands });
} }
default: { default:
const unknownCommand = command as { type: string }; return assertNeverRpcCommand(command);
return error(id, unknownCommand.type, `Unknown command: ${unknownCommand.type}`);
}
} }
} }