codekingpro/portable-devtools
114k
1"""In-process handle for a single tool call's streaming execution.2 3Mirrors the shape of `ChatModelStream` from langchain-core but simpler —4a tool has one output channel, no content-block multiplexing. Populated5by `ToolCallTransformer` as `tool-started` / `tool-output-delta` /6`tool-finished` / `tool-error` events flow in on the `tools` channel.7"""8 9from __future__ import annotations10 11from collections.abc import AsyncIterator, Iterator12from typing import Any13 14from langgraph.stream.stream_channel import StreamChannel15 16 17class ToolCallStream:18 """Scoped view of a single tool call's lifecycle.19 20 Yielded on `run.tool_calls` once per `tool-started` event. Fields21 are populated as events arrive:22 23 - `tool_call_id`, `tool_name`, `input`: stable from the start event.24 - `output_deltas`: a `StreamChannel` of delta chunks. Iterate (sync or25 async) to consume partial output in arrival order.26 - `output`: terminal payload from `tool-finished`, or `None` if the27 call failed or is still in flight.28 - `error`: terminal error string from `tool-error`, or `None` if the29 call succeeded or is still in flight.30 - `completed`: True once a terminal event (`tool-finished` or31 `tool-error`) has been observed.32 33 `ToolCallStream` is not meant to be constructed by end users — it's34 produced by `ToolCallTransformer` as events flow through the mux.35 """36 37 def __init__(38 self,39 tool_call_id: str,40 tool_name: str,41 input: dict[str, Any] | None = None,42 ) -> None:43 """Initialize a fresh handle for a tool call.44 45 Args:46 tool_call_id: The `tool_call_id` from the AIMessage.47 tool_name: The tool's name.48 input: The tool's input arguments (as reported by49 `on_tool_start`), or `None` if none were captured.50 """51 self.tool_call_id = tool_call_id52 self.tool_name = tool_name53 self.input = input54 self._output_deltas: StreamChannel[Any] = StreamChannel()55 self.output: Any = None56 self.error: str | None = None57 self.completed = False58 59 @property60 def output_deltas(self) -> StreamChannel[Any]:61 """The channel of streamed `tool-output-delta` payloads.62 63 Iterate (sync or async depending on how the run was started)64 to consume partial output in arrival order. The log closes when65 the tool finishes or errors.66 """67 return self._output_deltas68 69 def _bind(self, *, is_async: bool) -> None:70 """Bind the deltas log to sync or async iteration.71 72 Called by `ToolCallTransformer` when constructing this handle so73 the log matches the enclosing mux's mode.74 """75 self._output_deltas._bind(is_async=is_async)76 77 def _push_delta(self, delta: Any) -> None:78 self._output_deltas.push(delta)79 80 def _finish(self, output: Any) -> None:81 self.output = output82 self.completed = True83 self._output_deltas.close()84 85 def _fail(self, message: str) -> None:86 self.error = message87 self.completed = True88 self._output_deltas.close()89 90 def __iter__(self) -> Iterator[Any]:91 """Iterate delta chunks synchronously.92 93 Equivalent to `iter(self.output_deltas)`. Raises `TypeError` if94 the underlying log is bound to async mode.95 """96 return iter(self._output_deltas)97 98 def __aiter__(self) -> AsyncIterator[Any]:99 """Iterate delta chunks asynchronously.100 101 Equivalent to `aiter(self.output_deltas)`. Raises `TypeError`102 if the underlying log is bound to sync mode.103 """104 return self._output_deltas.__aiter__()105 106 def __repr__(self) -> str:107 status = (108 "completed"109 if self.completed and self.error is None110 else "failed"111 if self.completed112 else "running"113 )114 return (115 f"ToolCallStream(tool_call_id={self.tool_call_id!r}, "116 f"tool_name={self.tool_name!r}, status={status})"117 )118 