(env: Envelope)
| 57 | const env = result.value; |
| 58 | if (env.kind === Envelope_Kind.EVENT && env.payload.case === "askEvent") { |
| 59 | let payload: unknown = env.payload.value.payloadJson; |
| 60 | try { |
| 61 | payload = JSON.parse(env.payload.value.payloadJson); |
| 62 | } catch { |
| 63 | // Preserve a non-JSON payload verbatim. |
| 64 | } |
| 65 | yield { |
| 66 | type: "event", |
| 67 | eventType: env.payload.value.eventType, |
| 68 | payload, |
| 69 | machineId: env.payload.value.machineId, |
| 70 | }; |
| 71 | } else if (env.kind === Envelope_Kind.EVENT && env.payload.case === "askStreamError") { |
| 72 | yield { |
| 73 | type: "error", |
| 74 | code: STREAM_ERROR_NAMES[env.payload.value.code] ?? "internal", |
| 75 | message: env.payload.value.message, |
| 76 | }; |
| 77 | } else if (env.kind === Envelope_Kind.DONE && env.payload.case === "askDone") { |
| 78 | const value = env.payload.value; |
| 79 | yield { type: "done", ...value }; |
| 80 | } |
| 81 | } |
| 82 | } |
| 83 | |
| 84 | async collect(): Promise<string> { |
| 85 | let text = ""; |
| 86 | let streamError: Extract<AskFrame, { type: "error" }> | undefined; |
| 87 | for await (const frame of this) { |
| 88 | if (frame.type === "event") { |
| 89 | const payload = frame.payload as { delta?: string; text?: string } | undefined; |
| 90 | text += payload?.delta ?? payload?.text ?? ""; |
| 91 | } else if (frame.type === "error") { |
| 92 | streamError = frame; |
| 93 | } else if (!frame.success) { |
| 94 | throw new AgentStreamError( |
| 95 | frame.error ?? streamError?.message ?? "agent run failed", |
| 96 | streamError?.code ?? "internal", |
| 97 | frame.upstreamStatus, |
| 98 | ); |
| 99 | } else { |
| 100 | return frame.text || text; |
| 101 | } |
| 102 | } |
| 103 | return text; |
| 104 | } |
| 105 | } |
| 106 | |
| 107 | async function withTimeout<T>(promise: Promise<T>, timeoutMs: number): Promise<T> { |
| 108 | let timer: ReturnType<typeof setTimeout> | undefined; |
| 109 | try { |
| 110 | return await Promise.race([ |
| 111 | promise, |
| 112 | new Promise<never>((_, reject) => { |
no outgoing calls
no test coverage detected