Update files

This commit is contained in:
Mohamed Boudra
2026-02-04 09:29:26 +07:00
parent 7c15b26329
commit 10a2ef34a3
7 changed files with 640 additions and 327 deletions

View File

@@ -171,6 +171,8 @@ export function AgentList({
queryKey,
queryFn: async () => await client.getCheckoutStatus(agent.cwd),
staleTime: CHECKOUT_STATUS_STALE_TIME,
}).catch((error) => {
console.warn("[checkout_status] prefetch failed", error);
});
}
}

View File

@@ -311,6 +311,8 @@ export function GroupedAgentList({
queryKey,
queryFn: async () => await client.getCheckoutStatus(agent.cwd),
staleTime: CHECKOUT_STATUS_STALE_TIME,
}).catch((error) => {
console.warn("[checkout_status] prefetch failed", error);
});
}
}, [agents, queryClient]);
@@ -411,7 +413,9 @@ export function GroupedAgentList({
const session = useSessionStore.getState().sessions[agent.serverId];
const client = session?.client ?? null;
if (client) {
void client.archiveAgent(agent.id);
void client.archiveAgent(agent.id).catch((error) => {
console.warn("[archive_agent] failed", error);
});
}
},
[]

View File

@@ -137,4 +137,38 @@ describe("DaemonClientV2", () => {
isGit: false,
});
});
test("cancels waiters when send fails (no leaked timeouts)", async () => {
vi.useFakeTimers();
const logger = createMockLogger();
const mock = createMockTransport();
const transportFactory = () => ({
...mock.transport,
send: () => {
throw new Error("boom");
},
});
const client = new DaemonClientV2({
url: "ws://test",
logger,
reconnect: { enabled: false },
transportFactory,
});
clients.push(client);
const connectPromise = client.connect();
mock.triggerOpen();
await connectPromise;
const promise = client.getCheckoutStatus("/tmp/project");
await expect(promise).rejects.toThrow("boom");
// Ensure we didn't leave a waiter behind that will reject later.
expect((client as any).waiters.size).toBe(0);
vi.runOnlyPendingTimers();
vi.useRealTimers();
});
});

File diff suppressed because it is too large Load Diff

View File

@@ -1015,17 +1015,36 @@ export class Session {
break;
}
} catch (error: any) {
const err = error instanceof Error ? error : new Error(String(error));
this.sessionLogger.error(
{ err: error },
{ err },
"Error handling message"
);
const requestId = (msg as { requestId?: unknown }).requestId;
if (typeof requestId === "string") {
try {
this.emit({
type: "rpc_error",
payload: {
requestId,
requestType: msg.type,
error: "Request failed",
code: "handler_error",
},
});
} catch (emitError) {
this.sessionLogger.error({ err: emitError }, "Failed to emit rpc_error");
}
}
this.emit({
type: "activity_log",
payload: {
id: uuidv4(),
timestamp: new Date(),
type: "error",
content: `Error: ${error.message}`,
content: `Error: ${err.message}`,
},
});
}

View File

@@ -147,7 +147,56 @@ export class WebSocketSessionBridge {
private async handleRawMessage(ws: WebSocket, data: Buffer | ArrayBuffer | Buffer[]): Promise<void> {
try {
const parsed = JSON.parse(data.toString());
const message = WSInboundMessageSchema.parse(parsed);
const parsedMessage = WSInboundMessageSchema.safeParse(parsed);
if (!parsedMessage.success) {
const requestInfo = extractRequestInfoFromUnknownWsInbound(parsed);
const isUnknownSchema =
requestInfo?.requestId != null &&
typeof parsed === "object" &&
parsed != null &&
"type" in parsed &&
(parsed as { type?: unknown }).type === "session";
this.logger.warn(
{
requestId: requestInfo?.requestId,
requestType: requestInfo?.requestType,
error: parsedMessage.error.message,
},
"WS inbound message validation failed"
);
if (requestInfo) {
this.sendToClient(
ws,
wrapSessionMessage({
type: "rpc_error",
payload: {
requestId: requestInfo.requestId,
requestType: requestInfo.requestType,
error: isUnknownSchema ? "Unknown request schema" : "Invalid message",
code: isUnknownSchema ? "unknown_schema" : "invalid_message",
},
})
);
return;
}
const errorMessage = `Invalid message: ${parsedMessage.error.message}`;
this.sendToClient(
ws,
wrapSessionMessage({
type: "status",
payload: {
status: "error",
message: errorMessage,
},
})
);
return;
}
const message = parsedMessage.data;
const messageSummary = {
type: message.type,
@@ -224,6 +273,23 @@ export class WebSocketSessionBridge {
"Failed to parse/handle message"
);
const requestInfo = extractRequestInfoFromUnknownWsInbound(parsedPayload);
if (requestInfo) {
this.sendToClient(
ws,
wrapSessionMessage({
type: "rpc_error",
payload: {
requestId: requestInfo.requestId,
requestType: requestInfo.requestType,
error: "Invalid message",
code: "invalid_message",
},
})
);
return;
}
this.sendToClient(
ws,
wrapSessionMessage({
@@ -456,3 +522,38 @@ export class WebSocketSessionBridge {
}
}
}
function extractRequestInfoFromUnknownWsInbound(
payload: unknown
): { requestId: string; requestType?: string } | null {
if (!payload || typeof payload !== "object") {
return null;
}
const record = payload as {
type?: unknown;
requestId?: unknown;
message?: unknown;
};
// Session-wrapped messages
if (record.type === "session" && record.message && typeof record.message === "object") {
const msg = record.message as { requestId?: unknown; type?: unknown };
if (typeof msg.requestId === "string") {
return {
requestId: msg.requestId,
...(typeof msg.type === "string" ? { requestType: msg.type } : {}),
};
}
}
// Non-session messages (future-proof)
if (typeof record.requestId === "string") {
return {
requestId: record.requestId,
...(typeof record.type === "string" ? { requestType: record.type } : {}),
};
}
return null;
}

View File

@@ -1024,6 +1024,16 @@ export const StatusMessageSchema = z.object({
.passthrough(), // Allow additional fields
});
export const RpcErrorMessageSchema = z.object({
type: z.literal("rpc_error"),
payload: z.object({
requestId: z.string(),
requestType: z.string().optional(),
error: z.string(),
code: z.string().optional(),
}),
});
const AgentStatusWithRequestSchema = z.object({
agentId: z.string(),
requestId: z.string(),
@@ -1607,6 +1617,7 @@ export const SessionOutboundMessageSchema = z.discriminatedUnion("type", [
DictationStreamFinalMessageSchema,
DictationStreamErrorMessageSchema,
StatusMessageSchema,
RpcErrorMessageSchema,
InitializeAgentResponseMessageSchema,
ArtifactMessageSchema,
VoiceConversationLoadedMessageSchema,
@@ -1663,6 +1674,7 @@ export type AssistantChunkMessage = z.infer<typeof AssistantChunkMessageSchema>;
export type AudioOutputMessage = z.infer<typeof AudioOutputMessageSchema>;
export type TranscriptionResultMessage = z.infer<typeof TranscriptionResultMessageSchema>;
export type StatusMessage = z.infer<typeof StatusMessageSchema>;
export type RpcErrorMessage = z.infer<typeof RpcErrorMessageSchema>;
export type ArtifactMessage = z.infer<typeof ArtifactMessageSchema>;
export type VoiceConversationLoadedMessage = z.infer<
typeof VoiceConversationLoadedMessageSchema