codekingpro/portable-devtools
114k
1from __future__ import annotations2 3from abc import abstractmethod4from collections.abc import AsyncIterator, Callable, Iterator, Sequence5from typing import Any, Generic, Literal, cast, overload6 7from langchain_core.runnables import Runnable, RunnableConfig8from langchain_core.runnables.graph import Graph as DrawableGraph9from typing_extensions import Self10 11from langgraph.types import (12 All,13 Command,14 GraphOutput,15 StateSnapshot,16 StateUpdate,17 StreamMode,18 StreamPart,19)20from langgraph.typing import ContextT, InputT, OutputT, StateT21 22__all__ = ("PregelProtocol", "StreamProtocol")23 24 25class PregelProtocol(Runnable[InputT, Any], Generic[StateT, ContextT, InputT, OutputT]):26 @abstractmethod27 def with_config(28 self, config: RunnableConfig | None = None, **kwargs: Any29 ) -> Self: ...30 31 @abstractmethod32 def get_graph(33 self,34 config: RunnableConfig | None = None,35 *,36 xray: int | bool = False,37 ) -> DrawableGraph: ...38 39 @abstractmethod40 async def aget_graph(41 self,42 config: RunnableConfig | None = None,43 *,44 xray: int | bool = False,45 ) -> DrawableGraph: ...46 47 @abstractmethod48 def get_state(49 self, config: RunnableConfig, *, subgraphs: bool = False50 ) -> StateSnapshot: ...51 52 @abstractmethod53 async def aget_state(54 self, config: RunnableConfig, *, subgraphs: bool = False55 ) -> StateSnapshot: ...56 57 @abstractmethod58 def get_state_history(59 self,60 config: RunnableConfig,61 *,62 filter: dict[str, Any] | None = None,63 before: RunnableConfig | None = None,64 limit: int | None = None,65 ) -> Iterator[StateSnapshot]: ...66 67 @abstractmethod68 def aget_state_history(69 self,70 config: RunnableConfig,71 *,72 filter: dict[str, Any] | None = None,73 before: RunnableConfig | None = None,74 limit: int | None = None,75 ) -> AsyncIterator[StateSnapshot]: ...76 77 @abstractmethod78 def bulk_update_state(79 self,80 config: RunnableConfig,81 updates: Sequence[Sequence[StateUpdate]],82 ) -> RunnableConfig: ...83 84 @abstractmethod85 async def abulk_update_state(86 self,87 config: RunnableConfig,88 updates: Sequence[Sequence[StateUpdate]],89 ) -> RunnableConfig: ...90 91 @abstractmethod92 def update_state(93 self,94 config: RunnableConfig,95 values: dict[str, Any] | Any | None,96 as_node: str | None = None,97 ) -> RunnableConfig: ...98 99 @abstractmethod100 async def aupdate_state(101 self,102 config: RunnableConfig,103 values: dict[str, Any] | Any | None,104 as_node: str | None = None,105 ) -> RunnableConfig: ...106 107 @overload108 @abstractmethod109 def stream(110 self,111 input: InputT | Command | None,112 config: RunnableConfig | None = None,113 *,114 context: ContextT | None = None,115 stream_mode: StreamMode | list[StreamMode] | None = None,116 interrupt_before: All | Sequence[str] | None = None,117 interrupt_after: All | Sequence[str] | None = None,118 subgraphs: bool = False,119 version: Literal["v2"],120 ) -> Iterator[StreamPart[StateT, OutputT]]: ...121 122 @overload123 @abstractmethod124 def stream(125 self,126 input: InputT | Command | None,127 config: RunnableConfig | None = None,128 *,129 context: ContextT | None = None,130 stream_mode: StreamMode | list[StreamMode] | None = None,131 interrupt_before: All | Sequence[str] | None = None,132 interrupt_after: All | Sequence[str] | None = None,133 subgraphs: bool = False,134 version: Literal["v1"] = ...,135 ) -> Iterator[dict[str, Any] | Any]: ...136 137 @abstractmethod138 def stream(139 self,140 input: InputT | Command | None,141 config: RunnableConfig | None = None,142 *,143 context: ContextT | None = None,144 stream_mode: StreamMode | list[StreamMode] | None = None,145 interrupt_before: All | Sequence[str] | None = None,146 interrupt_after: All | Sequence[str] | None = None,147 subgraphs: bool = False,148 version: Literal["v1", "v2"] = "v1",149 ) -> Iterator[dict[str, Any] | Any]: ...150 151 @overload152 @abstractmethod153 def astream(154 self,155 input: InputT | Command | None,156 config: RunnableConfig | None = None,157 *,158 context: ContextT | None = None,159 stream_mode: StreamMode | list[StreamMode] | None = None,160 interrupt_before: All | Sequence[str] | None = None,161 interrupt_after: All | Sequence[str] | None = None,162 subgraphs: bool = False,163 version: Literal["v2"],164 ) -> AsyncIterator[StreamPart[StateT, OutputT]]: ...165 166 @overload167 @abstractmethod168 def astream(169 self,170 input: InputT | Command | None,171 config: RunnableConfig | None = None,172 *,173 context: ContextT | None = None,174 stream_mode: StreamMode | list[StreamMode] | None = None,175 interrupt_before: All | Sequence[str] | None = None,176 interrupt_after: All | Sequence[str] | None = None,177 subgraphs: bool = False,178 version: Literal["v1"] = ...,179 ) -> AsyncIterator[dict[str, Any] | Any]: ...180 181 @abstractmethod182 def astream(183 self,184 input: InputT | Command | None,185 config: RunnableConfig | None = None,186 *,187 context: ContextT | None = None,188 stream_mode: StreamMode | list[StreamMode] | None = None,189 interrupt_before: All | Sequence[str] | None = None,190 interrupt_after: All | Sequence[str] | None = None,191 subgraphs: bool = False,192 version: Literal["v1", "v2"] = "v1",193 ) -> AsyncIterator[dict[str, Any] | Any]: ...194 195 @overload196 @abstractmethod197 def invoke(198 self,199 input: InputT | Command | None,200 config: RunnableConfig | None = None,201 *,202 context: ContextT | None = None,203 interrupt_before: All | Sequence[str] | None = None,204 interrupt_after: All | Sequence[str] | None = None,205 version: Literal["v2"],206 ) -> GraphOutput[OutputT]: ...207 208 @overload209 @abstractmethod210 def invoke(211 self,212 input: InputT | Command | None,213 config: RunnableConfig | None = None,214 *,215 context: ContextT | None = None,216 interrupt_before: All | Sequence[str] | None = None,217 interrupt_after: All | Sequence[str] | None = None,218 version: Literal["v1"] = ...,219 ) -> dict[str, Any] | Any: ...220 221 @abstractmethod222 def invoke(223 self,224 input: InputT | Command | None,225 config: RunnableConfig | None = None,226 *,227 context: ContextT | None = None,228 interrupt_before: All | Sequence[str] | None = None,229 interrupt_after: All | Sequence[str] | None = None,230 version: Literal["v1", "v2"] = "v1",231 ) -> dict[str, Any] | Any: ...232 233 @overload234 @abstractmethod235 async def ainvoke(236 self,237 input: InputT | Command | None,238 config: RunnableConfig | None = None,239 *,240 context: ContextT | None = None,241 interrupt_before: All | Sequence[str] | None = None,242 interrupt_after: All | Sequence[str] | None = None,243 version: Literal["v2"],244 ) -> GraphOutput[OutputT]: ...245 246 @overload247 @abstractmethod248 async def ainvoke(249 self,250 input: InputT | Command | None,251 config: RunnableConfig | None = None,252 *,253 context: ContextT | None = None,254 interrupt_before: All | Sequence[str] | None = None,255 interrupt_after: All | Sequence[str] | None = None,256 version: Literal["v1"] = ...,257 ) -> dict[str, Any] | Any: ...258 259 @abstractmethod260 async def ainvoke(261 self,262 input: InputT | Command | None,263 config: RunnableConfig | None = None,264 *,265 context: ContextT | None = None,266 interrupt_before: All | Sequence[str] | None = None,267 interrupt_after: All | Sequence[str] | None = None,268 version: Literal["v1", "v2"] = "v1",269 ) -> dict[str, Any] | Any: ...270 271 272StreamChunk = tuple[tuple[str, ...], str, Any]273 274 275class StreamProtocol:276 __slots__ = ("modes", "__call__")277 278 modes: set[StreamMode]279 280 __call__: Callable[[Self, StreamChunk], None]281 282 def __init__(283 self,284 __call__: Callable[[StreamChunk], None],285 modes: set[StreamMode],286 ) -> None:287 self.__call__ = cast(Callable[[Self, StreamChunk], None], __call__)288 self.modes = modes289 