codekingpro/portable-devtools
114k
1from __future__ import annotations2 3from collections.abc import Sequence4from typing import Any, Generic5 6from typing_extensions import Self7 8from langgraph._internal._typing import MISSING9from langgraph.channels.base import BaseChannel, Value10from langgraph.errors import (11 EmptyChannelError,12 ErrorCode,13 InvalidUpdateError,14 create_error_message,15)16 17__all__ = ("LastValue", "LastValueAfterFinish")18 19 20class LastValue(Generic[Value], BaseChannel[Value, Value, Value]):21 """Stores the last value received, can receive at most one value per step."""22 23 __slots__ = ("value",)24 25 value: Value | Any26 27 def __init__(self, typ: Any, key: str = "") -> None:28 super().__init__(typ, key)29 self.value = MISSING30 31 def __eq__(self, value: object) -> bool:32 return isinstance(value, LastValue)33 34 @property35 def ValueType(self) -> type[Value]:36 """The type of the value stored in the channel."""37 return self.typ38 39 @property40 def UpdateType(self) -> type[Value]:41 """The type of the update received by the channel."""42 return self.typ43 44 def copy(self) -> Self:45 """Return a copy of the channel."""46 empty = self.__class__(self.typ, self.key)47 empty.value = self.value48 return empty49 50 def from_checkpoint(self, checkpoint: Value) -> Self:51 empty = self.__class__(self.typ, self.key)52 if checkpoint is not MISSING:53 empty.value = checkpoint54 return empty55 56 def update(self, values: Sequence[Value]) -> bool:57 if len(values) == 0:58 return False59 if len(values) != 1:60 msg = create_error_message(61 message=f"At key '{self.key}': Can receive only one value per step. Use an Annotated key to handle multiple values.",62 error_code=ErrorCode.INVALID_CONCURRENT_GRAPH_UPDATE,63 )64 raise InvalidUpdateError(msg)65 66 self.value = values[-1]67 return True68 69 def get(self) -> Value:70 if self.value is MISSING:71 raise EmptyChannelError()72 return self.value73 74 def is_available(self) -> bool:75 return self.value is not MISSING76 77 def checkpoint(self) -> Value:78 return self.value79 80 81class LastValueAfterFinish(82 Generic[Value], BaseChannel[Value, Value, tuple[Value, bool]]83):84 """Stores the last value received, but only made available after finish().85 Once made available, clears the value."""86 87 __slots__ = ("value", "finished")88 89 value: Value | Any90 finished: bool91 92 def __init__(self, typ: Any, key: str = "") -> None:93 super().__init__(typ, key)94 self.value = MISSING95 self.finished = False96 97 def __eq__(self, value: object) -> bool:98 return isinstance(value, LastValueAfterFinish)99 100 @property101 def ValueType(self) -> type[Value]:102 """The type of the value stored in the channel."""103 return self.typ104 105 @property106 def UpdateType(self) -> type[Value]:107 """The type of the update received by the channel."""108 return self.typ109 110 def checkpoint(self) -> tuple[Value | Any, bool] | Any:111 if self.value is MISSING:112 return MISSING113 return (self.value, self.finished)114 115 def from_checkpoint(self, checkpoint: tuple[Value | Any, bool] | Any) -> Self:116 empty = self.__class__(self.typ)117 empty.key = self.key118 if checkpoint is not MISSING:119 empty.value, empty.finished = checkpoint120 return empty121 122 def update(self, values: Sequence[Value | Any]) -> bool:123 if len(values) == 0:124 return False125 126 self.finished = False127 self.value = values[-1]128 return True129 130 def consume(self) -> bool:131 if self.finished:132 self.finished = False133 self.value = MISSING134 return True135 136 return False137 138 def finish(self) -> bool:139 if not self.finished and self.value is not MISSING:140 self.finished = True141 return True142 else:143 return False144 145 def get(self) -> Value:146 if self.value is MISSING or not self.finished:147 raise EmptyChannelError()148 return self.value149 150 def is_available(self) -> bool:151 return self.value is not MISSING and self.finished152 