codekingpro/portable-devtools
115k
1from __future__ import annotations2 3from collections.abc import Callable4from dataclasses import dataclass, field, replace5from typing import Any, Generic, cast6 7from langgraph.store.base import BaseStore8from langgraph_sdk.auth.types import BaseUser9from typing_extensions import TypedDict, Unpack10 11from langgraph._internal._constants import CONF, CONFIG_KEY_RUNTIME12from langgraph.config import get_config13from langgraph.types import _DC_KWARGS, StreamWriter14from langgraph.typing import ContextT15 16__all__ = (17 "BaseUser",18 "ExecutionInfo",19 "RunControl",20 "Runtime",21 "ServerInfo",22 "get_runtime",23)24 25 26@dataclass(frozen=True, slots=True)27class ExecutionInfo:28 """Read-only execution info/metadata for the execution of current thread/run/node."""29 30 checkpoint_id: str31 """The checkpoint ID for the current execution."""32 33 checkpoint_ns: str34 """The checkpoint namespace for the current execution."""35 36 task_id: str37 """The task ID for the current execution."""38 39 thread_id: str | None = None40 """The thread ID for the current execution.41 42 None when running without a checkpointer (i.e., no persistence)."""43 44 run_id: str | None = None45 """The run ID for the current execution.46 47 None when `run_id` is not provided in the RunnableConfig."""48 49 node_attempt: int = 150 """Current node execution attempt number (1-indexed)."""51 52 node_first_attempt_time: float | None = None53 """Unix timestamp (seconds) for when the first attempt started."""54 55 def patch(self, **overrides: Any) -> ExecutionInfo:56 """Return a new execution info object with selected fields replaced."""57 return replace(self, **overrides)58 59 60@dataclass(frozen=True, slots=True)61class ServerInfo:62 """Metadata injected by LangGraph Server. None when running open-source LangGraph without LangSmith deployments."""63 64 assistant_id: str65 """The assistant ID for the current execution."""66 67 graph_id: str68 """The graph ID for the current execution."""69 70 user: BaseUser | None = None71 """The authenticated user, if any.72 73 This implements the `BaseUser` protocol from `langgraph_sdk.auth.types`,74 which supports both attribute access (e.g. `user.identity`) and dict-like75 access (e.g. `user["identity"]`).76 """77 78 79class RunControl:80 """Run-scoped control surface for cooperative draining.81 82 Intended for a single graph run. Create a fresh `RunControl` per run;83 reusing a control after `request_drain()` leaves it drained.84 85 Safe to call from any thread: the drain request is represented by a86 single attribute write, so no lock is needed for this signal.87 If more mutable state is added here, add synchronization.88 """89 90 __slots__ = ("_drain_reason",)91 92 def __init__(self) -> None:93 self._drain_reason: str | None = None94 95 def request_drain(self, reason: str = "shutdown") -> None:96 self._drain_reason = reason97 98 @property99 def drain_requested(self) -> bool:100 return self._drain_reason is not None101 102 @property103 def drain_reason(self) -> str | None:104 return self._drain_reason105 106 107def _no_op_stream_writer(_: Any) -> None: ...108 109 110def _no_op_heartbeat() -> None: ...111 112 113class _RuntimeOverrides(TypedDict, Generic[ContextT], total=False):114 context: ContextT115 store: BaseStore | None116 stream_writer: StreamWriter117 heartbeat: Callable[[], None]118 previous: Any119 execution_info: ExecutionInfo120 server_info: ServerInfo | None121 control: RunControl | None122 123 124@dataclass(**_DC_KWARGS)125class Runtime(Generic[ContextT]):126 """Convenience class that bundles run-scoped context and other runtime utilities.127 128 This class is injected into graph nodes and middleware. It provides access to129 `context`, `store`, `stream_writer`, `previous`, and `execution_info`.130 131 !!! note "Accessing `config`"132 133 `Runtime` does not include `config`. To access `RunnableConfig`, you can inject134 it directly by adding a `config: RunnableConfig` parameter to your node function135 (recommended), or use `get_config()` from `langgraph.config`.136 137 !!! note138 `ToolRuntime` (from `langgraph.prebuilt`) is a subclass that provides similar139 functionality but is designed specifically for tools. It shares `context`, `store`,140 and `stream_writer` with `Runtime`, and adds tool-specific attributes like `config`,141 `state`, and `tool_call_id`.142 143 !!! version-added "Added in version v0.6.0"144 145 Example:146 147 ```python148 from typing import TypedDict149 from langgraph.graph import StateGraph150 from dataclasses import dataclass151 from langgraph.runtime import Runtime152 from langgraph.store.memory import InMemoryStore153 154 155 @dataclass156 class Context: # (1)!157 user_id: str158 159 160 class State(TypedDict, total=False):161 response: str162 163 164 store = InMemoryStore() # (2)!165 store.put(("users",), "user_123", {"name": "Alice"})166 167 168 def personalized_greeting(state: State, runtime: Runtime[Context]) -> State:169 '''Generate personalized greeting using runtime context and store.'''170 user_id = runtime.context.user_id # (3)!171 name = "unknown_user"172 if runtime.store:173 if memory := runtime.store.get(("users",), user_id):174 name = memory.value["name"]175 176 response = f"Hello {name}! Nice to see you again."177 return {"response": response}178 179 180 graph = (181 StateGraph(state_schema=State, context_schema=Context)182 .add_node("personalized_greeting", personalized_greeting)183 .set_entry_point("personalized_greeting")184 .set_finish_point("personalized_greeting")185 .compile(store=store)186 )187 188 result = graph.invoke({}, context=Context(user_id="user_123"))189 print(result)190 # > {'response': 'Hello Alice! Nice to see you again.'}191 ```192 193 1. Define a schema for the runtime context.194 2. Create a store to persist memories and other information.195 3. Use the runtime context to access the `user_id`.196 """197 198 context: ContextT = field(default=None) # type: ignore[assignment]199 """Static context for the graph run, like `user_id`, `db_conn`, etc.200 201 Can also be thought of as 'run dependencies'."""202 203 store: BaseStore | None = field(default=None)204 """Store for the graph run, enabling persistence and memory."""205 206 stream_writer: StreamWriter = field(default=_no_op_stream_writer)207 """Function that writes to the custom stream."""208 209 heartbeat: Callable[[], None] = field(default=_no_op_heartbeat)210 """Record progress for the current node's `idle_timeout`.211 212 Call this from inside long-running work that does not naturally emit213 writes, stream chunks, child tasks, or LangChain callback events, to214 prevent the node from being treated as idle. It is also the only215 progress signal honored under `TimeoutPolicy(refresh_on="heartbeat")`.216 Outside an idle-timed attempt this is a no-op.217 """218 219 previous: Any = field(default=None)220 """The previous return value for the given thread.221 222 Only available with the functional API when a checkpointer is provided.223 """224 225 execution_info: ExecutionInfo | None = field(default=None)226 """Read-only execution information/metadata for the current node run.227 228 None before task preparation populates it."""229 230 server_info: ServerInfo | None = field(default=None)231 """Metadata injected by LangGraph Server. None when running open-source LangGraph without LangSmith deployments."""232 233 control: RunControl | None = field(default=None)234 """Run-scoped control plane for cooperative draining.235 236 Populated automatically during graph runs. None outside an active237 graph runtime.238 """239 240 def merge(self, other: Runtime[ContextT]) -> Runtime[ContextT]:241 """Merge two runtimes together.242 243 If a value is not provided in the other runtime, the value from the current runtime is used.244 """245 return Runtime(246 context=other.context or self.context,247 store=other.store or self.store,248 stream_writer=other.stream_writer249 if other.stream_writer is not _no_op_stream_writer250 else self.stream_writer,251 heartbeat=other.heartbeat252 if other.heartbeat is not _no_op_heartbeat253 else self.heartbeat,254 previous=self.previous if other.previous is None else other.previous,255 execution_info=other.execution_info or self.execution_info,256 server_info=other.server_info or self.server_info,257 control=other.control or self.control,258 )259 260 def override(261 self, **overrides: Unpack[_RuntimeOverrides[ContextT]]262 ) -> Runtime[ContextT]:263 """Replace the runtime with a new runtime with the given overrides."""264 return replace(self, **overrides)265 266 def patch_execution_info(self, **overrides: Any) -> Runtime[ContextT]:267 """Return a new runtime with selected execution_info fields replaced."""268 if self.execution_info is None:269 msg = "Cannot patch execution_info before it has been set"270 raise RuntimeError(msg)271 return replace(272 self,273 execution_info=self.execution_info.patch(**overrides),274 )275 276 @property277 def drain_requested(self) -> bool:278 return self.control.drain_requested if self.control is not None else False279 280 @property281 def drain_reason(self) -> str | None:282 return self.control.drain_reason if self.control is not None else None283 284 285DEFAULT_RUNTIME = Runtime(286 context=None,287 store=None,288 stream_writer=_no_op_stream_writer,289 heartbeat=_no_op_heartbeat,290 previous=None,291 execution_info=None,292 control=None,293)294 295 296def get_runtime(context_schema: type[ContextT] | None = None) -> Runtime[ContextT]:297 """Get the runtime for the current graph run.298 299 Args:300 context_schema: Optional schema used for type hinting the return type of the runtime.301 302 Returns:303 The runtime for the current graph run.304 """305 306 # TODO: in an ideal world, we would have a context manager for307 # the runtime that's independent of the config. this will follow308 # from the removal of the configurable packing309 runtime = cast(Runtime[ContextT], get_config()[CONF].get(CONFIG_KEY_RUNTIME))310 return runtime311 