Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
protocol.py289 linesDownload Raw Back to pregel
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 
codekingpro/portable-devtools · Team Ai