MCPcopy Create free account
hub / github.com/commandoperator/cmdop-sdk / frameFromEnvelope

Function frameFromEnvelope

node/src/streaming.ts:59–109  ·  view source on GitHub ↗
(env: Envelope)

Source from the content-addressed store, hash-verified

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
107async 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) => {

Callers 1

Calls

no outgoing calls

Tested by

no test coverage detected