Project one streamed Envelope onto a typed frame.
(env: pb.Envelope)
| 114 | ) |
| 115 | elif item.kind == pb.Envelope.KIND_DONE: |
| 116 | done = item.ask_done |
| 117 | yield DoneFrame( |
| 118 | "done", |
| 119 | done.success, |
| 120 | done.text, |
| 121 | _optional(done, "error"), |
| 122 | _optional(done, "upstream_status"), |
| 123 | _optional(done, "provider"), |
| 124 | _optional(done, "detail"), |
| 125 | _optional(done, "input_tokens"), |
| 126 | _optional(done, "output_tokens"), |
| 127 | _optional(done, "total_tokens"), |
| 128 | _optional(done, "cache_read_tokens"), |
| 129 | _optional(done, "cache_write_tokens"), |
| 130 | _optional(done, "cost_usd"), |
| 131 | _optional(done, "cost_source"), |
| 132 | _optional(done, "engine"), |
| 133 | _optional(done, "model"), |
| 134 | ) |
| 135 | finally: |
| 136 | self._transport._abandon_stream(call_id) |
| 137 | |
| 138 | async def collect(self) -> str: |
| 139 | text = "" |
| 140 | stream_error: StreamErrorFrame | None = None |
| 141 | async for frame in self: |
| 142 | if isinstance(frame, EventFrame) and isinstance(frame.payload, dict): |
| 143 | text += frame.payload.get("delta") or frame.payload.get("text") or "" |
| 144 | elif isinstance(frame, StreamErrorFrame): |
| 145 | stream_error = frame |
| 146 | elif isinstance(frame, DoneFrame): |
| 147 | if not frame.success: |
| 148 | raise AgentStreamError( |
| 149 | frame.error or (stream_error.message if stream_error else "agent run failed"), |
| 150 | code=stream_error.code if stream_error else "internal", |
| 151 | status=frame.upstream_status, |
| 152 | ) |
| 153 | return frame.text or text |
| 154 | return text |
no test coverage detected