LangChain SDK Runtime
class App(BaseApp)LangChain SDK for MongoDB Atlas Agent Engine.
Example:
from agent_engine_sdk_langgraph import App
from langchain_openai import ChatOpenAI
app = App(app_name="My Agent")
@app.tool()
def my_tool(query: str) -> str:
return "result"
@app.entrypoint
def build_agent():
llm = app.llm(ChatOpenAI(model="gpt-4"))
checkpointer = app.checkpointer()
tools = app.get_tools()
# Build LangGraph...
return graphdef __init__(app_name: str,
app_version: str = "1.0.0",
mongodb_uri: str | None = None,
database_name: str | None = None,
enable_tracing: bool | None = None,
traces_collection_name: str = "traces",
enable_memory: bool | None = None,
org_id: str | None = None)Initialize the App.
Arguments:
app_name- Application nameapp_version- Application versionmongodb_uri- MongoDB URIdatabase_name- Database nameenable_tracing- Deprecated. Tracing is always enabled.traces_collection_name- Collection name for trace storageenable_memory- Deprecated. Use agent.yaml features.memory instead.org_id- Deprecated and ignored. The org is taken from theORG_IDenvironment variable, which the platform injects. Passing this argument raises a DeprecationWarning and will be removed in a future release.
@property
def agent_config() -> RuntimeAgentConfigReturn the runtime SDK's parsed view of agent.yaml.
@property
def memory() -> MemoryReturn the unified Memory facade over app-bound adapters.
The return type is agent_engine_sdk_memory.Memory. Methods return
typed results (Create*Result, MemoryChunk, ContextResponse).
Identity is supplied via per-call args, bind, or ambient
execution contextvars.
def get_tool_definitions() -> list[ToolDefinition]Return sdk-core ToolDefinitions for all registered tools.
def tools() -> list[Any]Return wrapped, LangGraph-specific tools for ToolNode.
def tool(is_local: bool = True,
*,
provider_type: str | None = None,
scopes: list[str] | None = None,
network: list[str] | None = None,
timeout: int = 30,
redact_fields: list[str] | None = None,
response_format: Literal["content",
"content_and_artifact"] = "content")Register a tool function.
Creates a LangChain tool and registers the raw function on the runtime for Tool Pod execution.
Arguments:
is_local- If True, runs in AER; if False, runs in Tool Pod (default: True)provider_type- Credential provider type for delegated auth (e.g. "github")scopes- OAuth scopes requested for delegated authnetwork- Allowed network hosts (for tool pod)timeout- Execution timeout in secondsredact_fields- Fields to redact in logsresponse_format- LangChain tool response format."content"(default) treats the return value as the ToolMessage content;"content_and_artifact"expects a(content, artifact)two-tuple, whereartifactis data for downstream code that is not sent to the model. The artifact must be JSON-serializable so it survives the OE round trip. Note it is still recorded in the execution log and is not covered byredact_fields(which redacts inputs only), so do not place secrets in it.
Returns:
Decorator function
def entrypoint(fn: Callable[..., Any]) -> Callable[..., Any]Mark the graph builder function.
Arguments:
fn- Function that builds and returns a CompiledStateGraph
Returns:
The original function (unchanged)
def prepare_agent_input(fn: PrepareAgentInput) -> PrepareAgentInputRegister a hook that builds graph input from the invocation.
The hook receives the framework-neutral AgentInput (the opaque
caller payload) and RequestContext, and returns a value LangGraph
can ingest directly — a graph state dict or a Command. It lets an
agent author control how the request becomes the graph's starting
input (e.g. fold extra payload context into the prompt, seed a system
message, or populate custom state keys).
Only fresh invocations flow through the hook; HITL resume stays platform-managed. When no hook is registered, the default message-wrapping behavior is used, so existing agents are unaffected.
Example:
@app.prepare_agent_input
def prepare(input: AgentInput, ctx: RequestContext) -> GraphInput:
payload = input.payload
return {
"messages": [HumanMessage(content=payload["message"])],
"extra": payload.get("extra"),
}Arguments:
fn- Callable(AgentInput, RequestContext) -> GraphInput.
Returns:
The original function (unchanged), so it can be used as a decorator.
def resolve_thread_id(fn: ResolveThreadId) -> ResolveThreadIdRegister a hook that builds the LangGraph checkpoint thread_id.
Callers manage Atlas Agent Engine session_id (and authenticated user_id).
The agent owns how those map to the LangGraph checkpoint key. When
registered, the hook's return value is used verbatim on every
invocation — fresh and resume — with no workspace suffix appended.
When no hook is registered, the adapter derives
session_id:workspace_id as today.
Custom keys are invisible to Atlas Agent Engine session-history queries
(/query/sessions*), which still look up only the default
session/workspace-derived keys. Agents that bypass workspace
scoping also own collision isolation within the checkpoint
database (for example when sharing a project store with other
agents).
Example:
@app.resolve_thread_id
def resolve_thread_id(ctx: RequestContext) -> str:
actor = email_to_actor_id(ctx.user_id or "")
return f"{ctx.session_id}__{actor}"Arguments:
fn- Callable(RequestContext) -> str.
Returns:
The original function (unchanged), so it can be used as a decorator.
def output_parser(cls: type[OutputParser]) -> type[OutputParser]Register a custom output parser subclass on this app.
The class is validated with issubclass here, at decoration time, so
a non-OutputParser fails on import. The adapter stores the class
without instantiating it; because parse and on_stream_error are
abstract, a subclass missing either raises TypeError the moment it
is instantiated to stream (once the producer wires that up).
Example:
@app.output_parser
class BriefParser(LangGraphOutputParser):
stream_modes = ("messages", "values")
async def parse(self, item, ctx): ...
async def on_stream_error(self, ctx, error): ...Arguments:
cls- AnOutputParsersubclass.
Returns:
The original class (unchanged), so it can be used as a decorator.
def warm_up() -> boolBuild and cache the graph ahead of the first /execute call.
Never raises: a failure returns False so the AER records the
one-shot optimization as incomplete, while the existing lazy build
remains the fallback for the first real get_agent() call.
def get_agent(callbacks: list[Any] | None = None) -> AnyBuild and return a BaseAgent instance.
Calls the registered entrypoint function to build the graph, then wraps it in a LangGraphBaseAgent adapter.
Arguments:
callbacks- Optional list of BaseExecutionCallback instances for observability (e.g., NodeExecutionLogger from the AER). Wrapped in LangGraphCallbackAdapter before passing to the graph.
Returns:
LangGraphBaseAgent instance implementing BaseAgent protocol
Raises:
RuntimeError- If no entrypoint function has been registered
def run(**kwargs) -> NoneStart the agent service.
Registers the builder function and starts the runtime server. This is called by the platform entrypoint dispatcher.
Arguments:
**kwargs- Runtime startup options passed to the runtime.grpc_portandlog_levelare supported. Deprecatedhostandhttp_portvalues are still accepted for backwards compatibility but are ignored by agent-engine-runner-shared.
Raises:
RuntimeError- If no @app.entrypoint has been registered
def llm(llm: BaseChatModel, llm_id: str | None = None) -> BaseChatModelWrap a LangChain LLM for audited I/O through the Orchestration Engine.
For agents with a single LLM, call without llm_id::
llm = app.llm(ChatOpenAI(model="gpt-5.4"))
For agents with multiple LLMs, every call must supply a unique llm_id::
llm_a = app.llm(ChatOpenAI(model="gpt-5.4"), llm_id="primary") llm_b = app.llm(ChatOpenAI(model="gpt-5.4-mini"), llm_id="fast")
Unnamed calls register under the sentinel id "__default__"; calling
app.llm() twice without an id therefore raises the same duplicate-id
error as registering the same named id twice.
Arguments:
llm- LangChain LLM instance.llm_id- Unique identifier for this LLM. Optional when the agent uses a single LLM.
Returns:
SecureWrappedLLM in AER mode; the raw llm in TOOL mode.
Raises:
ValueError- Ifllm_idhas already been registered (including a second unnamed call which collides on"__default__").
def deep_agent(llm: BaseChatModel,
*,
tools: Sequence[Any] | None = None,
subagents: Sequence[Any] | None = None,
system_prompt: str | None = None,
middleware: Sequence[Any] = (),
checkpointer: Any = _UNSET,
store: Any = None,
skills: list[str] | None = None,
backend: Any = None) -> AnyCreate a deep agent graph pre-configured with Atlas Agent Engine secure routing.
Wraps create_agent_engine_deep_agent() with:
SecureWrappedLLMas the model (all LLM calls route through OE)- Configurable backend (defaults to
AgentEngineToolPodBackend) - Checkpointer sentinel resolution (
_UNSET->app.checkpointer(),None-> disable, instance -> use directly) - SubAgent model validation (string model specs rejected to prevent OE bypass)
The returned CompiledStateGraph should be used inside an
@app.entrypoint function. app.get_agent() will then wrap it
in LangGraphBaseAgent for AER compatibility.
Arguments:
-
llm- Base LLM instance to wrap with SecureWrappedLLM. -
tools- Additional tools for the deep agent (merged with built-in tools). -
subagents- SubAgent specs. Each spec with amodelfield must pass aBaseChatModelinstance, not a string. -
system_prompt- Custom system instructions. -
middleware- Additional middleware, run after the SDK's default interrupt-recovery and durable-nesting middleware (see create_agent_engine_deep_agent). -
checkpointer- LangGraph checkpointer for state persistence. Defaults to_UNSETwhich resolves toapp.checkpointer()(MongoDB). PassNoneto disable checkpointing. Pass aBaseCheckpointSaverinstance to use a custom checkpointer. -
store- LangGraph store for skills and shared data. -
skills- Parent directories (asstr) for deepagents skill discovery. Each immediate child directory containingSKILL.mdis one skill. Each path must be astr—pathlib.Pathobjects are NOT accepted because deepagents' SkillsMiddleware calls.rstrip("/")on the value (usestr(path)at the call site if you start from aPath). Relative paths resolve against the directory containingagent.yaml;AGENTIC_SKILLS_DIRmay set a different skills root. WhenNone(default) no skills are loaded. See the "Skills" section below for runtime behavior. -
backend- Backend for filesystem/shell ops. Defaults toAgentEngineToolPodBackendwhenNone. Pass a custom backend (e.g.StoreBackend) to override. Security note: custom backends bypass the default OE-audited I/O path — use only in tests or with backends that provide equivalent auditing.Skills (progressive disclosure): The
skillsparameter enables deepagents' progressive-disclosure pattern — the LLM sees only skill metadata (name + description + path) on turn 1 and reads the full SKILL.md body on demand viaread_file(path). This keeps the system prompt small while still giving the agent access to deep domain knowledge.Directory layout and frontmatter: Each entry in
skills=is a parent source directory. At runtime,SkillsMiddlewarelists that directory through the configured backend and treats each immediate child containingSKILL.mdas one skill; discovery is not recursive::/app/skills/ code-review-style/ SKILL.md examples/ good.py security-checklist/ SKILL.md
deepagents validates each skill's frontmatter itself. It skips unreadable or unparsable frontmatter and skills missing
nameordescription; Agent Skills naming or directory-name violations are warned about but may still load. This SDK does not validate or filter declared paths.Subagent non-inheritance: Skills attached to the top-level agent DO NOT propagate to subagents. If a subagent needs the same skills, pass
skills=on its ownSubAgentspec. This is a deepagents-level behavior and may change in a future release.Runtime state: On the first turn,
SkillsMiddleware.before_agentloads the metadata intostate["skills_metadata"]and injects a## Skills Systemblock into the system prompt. After turn 1 the metadata is cached on the thread; edits to SKILL.md files on disk are NOT picked up inside the same thread (start a new thread to pick up edits).
Returns:
CompiledStateGraph ready for LangGraphBaseAgent wrapping.
Raises:
RuntimeError- If any SubAgent spec uses a string model.
def checkpointer() -> Any | NoneGet the request-scoped platform checkpointer for LangGraph.
The returned PlatformCheckpointer delegates per invocation from
trusted request context: native sessions use the existing
MongoDB-backed saver; durable_workflow sessions use attempt-local
scratch keyed by execution + attempt + fence. It is created lazily on
first call and cached. The underlying MongoClient is closed when
App.close() is called.
The database name is read from CHECKPOINT_DB_NAME when set (exact
override, no project scoping). Otherwise the existing
MDB_AGENTIC_STORE_DB / per-project store resolution is used
(default: "mdb_store").
Returns:
PlatformCheckpointer instance in AER mode, None otherwise
def close() -> NoneRelease resources held by this App instance.
Closes the MongoDB client opened by :meth:checkpointer, if any.
def get_tools() -> list[Any]Get wrapped tools for ToolNode.
In AER mode, returns tools wrapped with SecureToolWrapper that route all executions through OE for logging and policy enforcement. In other modes, returns the original LangChain tools.
Returns:
List of tools ready for ToolNode
def get_tool_schemas() -> list[Any]Get tool schemas for llm.bind_tools().
Returns:
List of original LangChain tools (not wrapped)
def suspend(reason: str, context: dict[str, Any]) -> strGenerates a suspend command. If a tool should suspend, return the result of this function.
def finish_session() -> SessionFinishStatusMark this session finished so the platform frees its compute now.
Call it when the agent is done with the session. The current turn keeps running and returns its result normally; once it completes, the platform cancels any live sub-agent runs and releases the session's AER and tool pods instead of holding them until the idle timeout expires.
Safe to call more than once: the first call returns REQUESTED, later ones ALREADY_REQUESTED. Outside an agent run (local scripts, tool pods) there is no session to finish and the call returns UNAVAILABLE without raising. Calling it after the turn has already ended - e.g. from background work scheduled during the turn but that runs after it - also returns UNAVAILABLE: by then nothing is listening for the request anymore, so reporting REQUESTED would promise a release that will never happen.
A turn that suspends for human review, or that fails, keeps its resources so it stays resumable and diagnosable; the session then falls back to the idle timeout.
def validate_llm_response(response: Any) -> AnyApply guardrails validation to LLM response.
Checks if the response has content (and no tool calls), runs guardrails validation, and returns a new AIMessage with validated content if modifications were made.
Arguments:
response- LangChain AIMessage or similar response object
Returns:
Original or modified response after validation
def validate_output(text: str) -> strValidate text output using configured guardrails.
Arguments:
text- Text to validate
Returns:
Validated (possibly modified) text
def a2a_tools() -> list[Any]Get LangChain tools for A2A agent discovery and invocation.
Returns StructuredTool instances that execute in the AER process (not via the tool pod) because A2A requires the execution context's OE URL and A2A JWT.
Bind them to the LLM and pass to ToolNode alongside any @app.tool tools::
a2a = app.a2a_tools() llm_with_tools = llm.bind_tools(app.get_tool_schemas() + a2a) ToolNode(app.get_tools() + a2a)
Returns:
List of LangChain tools for discover_agents and invoke_agent. Empty list when A2A is not enabled in agent.yaml.
def get_current_user_id() -> str | NoneGet current execution context user ID.
Returns:
User ID from the current execution context, or None
LangGraph BaseAgent adapter.
Wraps a LangGraph CompiledGraph to implement the framework-neutral BaseAgent
protocol from agent-engine-sdk. Native vs durable branching lives on
ExecutionSession; this module stores hooks and drives invoke/stream.
class InternalStreamError(Exception)Stand-in passed to on_stream_error in place of the real exception.
The real exception may be a platform-internal one (LangGraph, an LLM
client, network/auth libraries) whose message was never meant to reach
the customer-visible custom_event stream that on_stream_error's
return value feeds into. The real exception is still logged and re-raised
unchanged; only what the parser sees is replaced.
class LangGraphBaseAgent(BaseAgent)LangGraph adapter implementing the BaseAgent protocol.
Wraps a LangGraph CompiledGraph and provides the framework-neutral invoke/stream interface that the AER expects.
This adapter handles:
- Converting AgentInput to LangGraph state format (with HumanMessage)
- Converting Command(resume=...) for HITL resume flows
- Streaming with token-level events via StreamEvent
- Converting final results to AgentOutput
Resume state travels on the RequestContext, not the caller payload:
- ctx.resume: Boolean indicating resume vs fresh execution
- ctx.resume_data: Data for Command(resume=...)
- ctx.metadata["checkpoint_id"]: Checkpoint to resume from
Example:
from langgraph.graph import StateGraph
from agent_engine_sdk_langgraph import LangGraphBaseAgent
# Build your graph
graph = StateGraph(...)
compiled = graph.compile(checkpointer=...)
# Wrap it with optional callbacks
agent = LangGraphBaseAgent(compiled, callbacks=[my_callback])
# Use via BaseAgent protocol
output = await agent.execute(ctx, agent_input) # non-streaming
async for event in agent.execute(ctx, agent_input): # streaming
...def __init__(graph: CompiledStateGraph,
callbacks: list[Any] | None = None,
prepare_input: PrepareAgentInput | None = None,
output_parser: type[OutputParser] | None = None,
use_custom_parser: bool = False,
resolve_thread_id: ResolveThreadId | None = None,
durable_subgraphs: DurableSubgraphResolver | None = None) -> NoneInitialize the adapter with a compiled LangGraph.
Arguments:
graph- A compiled LangGraph StateGraph with checkpointercallbacks- Optional list of LangChain callbacks for observabilityprepare_input- Optional tenant hook (@app.prepare_agent_input) that builds the graph input for a fresh invocation. When None, the default message-wrapping is used.output_parser- OptionalOutputParsersubclass (@app.output_parser) used to shape custom stream output.use_custom_parser- When true,stream()dual-runsoutput_parserand emits its output ascustom_event, and Atlas Agent Engine also subscribes to LangGraph'scustomchannel soemit_custom_eventreaches the stream (including when the parser does not listcustom). Requires a registered parser at deploy /get_agenttime. When false the parser is inert and Atlas Agent Engine does not subscribe tocustom.resolve_thread_id- Optional tenant hook (@app.resolve_thread_id) that builds the LangGraph checkpointthread_id. When None, the adapter derivessession_id:workspace_id. When set, the return value is used verbatim (no workspace suffix). Appended after the established positional surface so existing positional callers keep bindingoutput_parser/use_custom_parser.durable_subgraphs- Resolver for directly composed compiled subgraphs. It is used only while an OE durable attempt is active.
def execute(ctx: RequestContext, input: AgentInput) -> ExecutionResultRun the agent. Await for a result, or iterate for streaming.
async def invoke(ctx: RequestContext, input: AgentInput) -> AgentOutputExecute the agent and return the final result.
async def resume(ctx: RequestContext, input: AgentInput,
decision: str) -> AgentOutputResume a suspended agent with a human decision.
Convenience method that sets up resume value fields and delegates to invoke().
resume is the HITL-continuation bit (Command vs a new HumanMessage,
skip branch copy). resume_data is the decision payload. Presence of
data cannot imply the flag: a resume with a missing payload must fail
closed, not look like a fresh turn.
async def stream(ctx: RequestContext,
input: AgentInput) -> AsyncIterator[StreamEvent]Stream agent execution with token-level events.
Handles both fresh executions and resume flows based on the context:
- ctx.resume: Boolean indicating resume vs fresh execution
- ctx.resume_data: Data for Command(resume=...)
- ctx.metadata["checkpoint_id"]: Checkpoint to resume from
Yields StreamEvent objects with event:
- "token": Streaming LLM token (data["content"], data["source"], data["tool_call_id"])
- "subagent_start": Subagent dispatch begins (data["source"], data["subagent_name"], data["tool_call_id"], data["description"])
- "subagent_end": Subagent dispatch completes (data["source"], data["subagent_name"], data["tool_call_id"], data["summary"] with the subagent's final response text; empty string for defensive drain ends)
- "result": Final completion (data with response + messages)
- "suspend": HITL interrupt (data with suspend_payload + opaque metadata)
On durable sessions the attempt's scratch is released when the stream ends.