Skip to content

Latest commit

 

History

History
922 lines (643 loc) · 28.5 KB

File metadata and controls

922 lines (643 loc) · 28.5 KB

Table of Contents

LangChain SDK Runtime

App

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 graph

__init__

def __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 name
  • app_version - Application version
  • mongodb_uri - MongoDB URI
  • database_name - Database name
  • enable_tracing - Deprecated. Tracing is always enabled.
  • traces_collection_name - Collection name for trace storage
  • enable_memory - Deprecated. Use agent.yaml features.memory instead.
  • org_id - Deprecated and ignored. The org is taken from the ORG_ID environment variable, which the platform injects. Passing this argument raises a DeprecationWarning and will be removed in a future release.

agent_config

@property
def agent_config() -> RuntimeAgentConfig

Return the runtime SDK's parsed view of agent.yaml.

memory

@property
def memory() -> Memory

Return 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.

get_tool_definitions

def get_tool_definitions() -> list[ToolDefinition]

Return sdk-core ToolDefinitions for all registered tools.

tools

def tools() -> list[Any]

Return wrapped, LangGraph-specific tools for ToolNode.

tool

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 auth
  • network - Allowed network hosts (for tool pod)
  • timeout - Execution timeout in seconds
  • redact_fields - Fields to redact in logs
  • response_format - LangChain tool response format. "content" (default) treats the return value as the ToolMessage content; "content_and_artifact" expects a (content, artifact) two-tuple, where artifact is 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 by redact_fields (which redacts inputs only), so do not place secrets in it.

Returns:

Decorator function

entrypoint

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)

prepare_agent_input

def prepare_agent_input(fn: PrepareAgentInput) -> PrepareAgentInput

Register 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.

resolve_thread_id

def resolve_thread_id(fn: ResolveThreadId) -> ResolveThreadId

Register 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.

output_parser

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 - An OutputParser subclass.

Returns:

The original class (unchanged), so it can be used as a decorator.

warm_up

def warm_up() -> bool

Build 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.

get_agent

def get_agent(callbacks: list[Any] | None = None) -> Any

Build 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

run

def run(**kwargs) -> None

Start 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_port and log_level are supported. Deprecated host and http_port values are still accepted for backwards compatibility but are ignored by agent-engine-runner-shared.

Raises:

  • RuntimeError - If no @app.entrypoint has been registered

llm

def llm(llm: BaseChatModel, llm_id: str | None = None) -> BaseChatModel

Wrap 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 - If llm_id has already been registered (including a second unnamed call which collides on "__default__").

deep_agent

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) -> Any

Create a deep agent graph pre-configured with Atlas Agent Engine secure routing.

Wraps create_agent_engine_deep_agent() with:

  • SecureWrappedLLM as 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 a model field must pass a BaseChatModel instance, 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 _UNSET which resolves to app.checkpointer() (MongoDB). Pass None to disable checkpointing. Pass a BaseCheckpointSaver instance to use a custom checkpointer.

  • store - LangGraph store for skills and shared data.

  • skills - Parent directories (as str) for deepagents skill discovery. Each immediate child directory containing SKILL.md is one skill. Each path must be a str — pathlib.Path objects are NOT accepted because deepagents' SkillsMiddleware calls .rstrip("/") on the value (use str(path) at the call site if you start from a Path). Relative paths resolve against the directory containing agent.yaml; AGENTIC_SKILLS_DIR may set a different skills root. When None (default) no skills are loaded. See the "Skills" section below for runtime behavior.

  • backend - Backend for filesystem/shell ops. Defaults to AgentEngineToolPodBackend when None. 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 skills parameter 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 via read_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, SkillsMiddleware lists that directory through the configured backend and treats each immediate child containing SKILL.md as 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 name or description; 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 own SubAgent spec. This is a deepagents-level behavior and may change in a future release.

    Runtime state: On the first turn, SkillsMiddleware.before_agent loads the metadata into state["skills_metadata"] and injects a ## Skills System block 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.

checkpointer

def checkpointer() -> Any | None

Get 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

close

def close() -> None

Release resources held by this App instance.

Closes the MongoDB client opened by :meth:checkpointer, if any.

get_tools

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

get_tool_schemas

def get_tool_schemas() -> list[Any]

Get tool schemas for llm.bind_tools().

Returns:

List of original LangChain tools (not wrapped)

suspend

def suspend(reason: str, context: dict[str, Any]) -> str

Generates a suspend command. If a tool should suspend, return the result of this function.

finish_session

def finish_session() -> SessionFinishStatus

Mark 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.

validate_llm_response

def validate_llm_response(response: Any) -> Any

Apply 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

validate_output

def validate_output(text: str) -> str

Validate text output using configured guardrails.

Arguments:

  • text - Text to validate

Returns:

Validated (possibly modified) text

a2a_tools

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.

get_current_user_id

def get_current_user_id() -> str | None

Get 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.

InternalStreamError

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.

LangGraphBaseAgent

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
        ...

__init__

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) -> None

Initialize the adapter with a compiled LangGraph.

Arguments:

  • graph - A compiled LangGraph StateGraph with checkpointer
  • callbacks - Optional list of LangChain callbacks for observability
  • prepare_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 - Optional OutputParser subclass (@app.output_parser) used to shape custom stream output.
  • use_custom_parser - When true, stream() dual-runs output_parser and emits its output as custom_event, and Atlas Agent Engine also subscribes to LangGraph's custom channel so emit_custom_event reaches the stream (including when the parser does not list custom). Requires a registered parser at deploy / get_agent time. When false the parser is inert and Atlas Agent Engine does not subscribe to custom.
  • resolve_thread_id - Optional tenant hook (@app.resolve_thread_id) that builds the LangGraph checkpoint thread_id. When None, the adapter derives session_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 binding output_parser / use_custom_parser.
  • durable_subgraphs - Resolver for directly composed compiled subgraphs. It is used only while an OE durable attempt is active.

execute

def execute(ctx: RequestContext, input: AgentInput) -> ExecutionResult

Run the agent. Await for a result, or iterate for streaming.

invoke

async def invoke(ctx: RequestContext, input: AgentInput) -> AgentOutput

Execute the agent and return the final result.

resume

async def resume(ctx: RequestContext, input: AgentInput,
                 decision: str) -> AgentOutput

Resume 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.

stream

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.