Metadata-Version: 2.5
Name: aictrl-langgraph
Version: 0.2.0
Summary: Graph and execution monitoring for LangGraph workflows
Project-URL: Repository, https://github.com/aictrl-azenohi/ai-control-plane-pro
Project-URL: Documentation, https://github.com/aictrl-azenohi/ai-control-plane-pro/tree/main/packages/aictrl-langgraph
Project-URL: Issues, https://github.com/aictrl-azenohi/ai-control-plane-pro/issues
License-Expression: Apache-2.0
License-File: LICENSE
Keywords: agents,langgraph,monitoring,observability
Classifier: Development Status :: 4 - Beta
Classifier: Programming Language :: Python :: 3 :: Only
Classifier: Programming Language :: Python :: 3.11
Classifier: Programming Language :: Python :: 3.12
Classifier: Programming Language :: Python :: 3.13
Classifier: Programming Language :: Python :: 3.14
Classifier: Typing :: Typed
Requires-Python: <3.15,>=3.11
Requires-Dist: httpx<1,>=0.28
Requires-Dist: langchain-core<2,>=1.6
Requires-Dist: langgraph<2,>=1.2
Description-Content-Type: text/markdown

# AI Control Plane · LangGraph SDK

Connect a compiled LangGraph `StateGraph` to graph and run monitoring with one
Python package. The SDK registers the topology, reports run and node lifecycle
changes, and keeps the application connected with a heartbeat.
Optional policy guards enforce assignments configured in the graph workspace.

Python 3.11–3.14. Supports LangGraph `>=1.2,<2`, LangChain Core `>=1.6,<2` and
HTTPX `>=0.28,<1`; CI tests the lowest and latest allowed versions. License: Apache-2.0.

## Install

```sh
uv add aictrl-langgraph        # or: pip install aictrl-langgraph
```

Python 3.11–3.14. Pin a compatible range in applications, for example
`aictrl-langgraph>=0.2,<0.3`, while the SDK is in `0.x`. Releases and notes:
https://github.com/aictrl-azenohi/ai-control-plane-pro/releases

## Install from a local build

From the repository root:

```sh
uv build --package aictrl-langgraph --out-dir dist/sdk
```

In the agent application's environment:

```sh
python -m pip install /path/to/dist/sdk/aictrl_langgraph-0.2.0-py3-none-any.whl
```

The wheel is standalone; it does not require the repository or the backend package.
Released versions are on PyPI (see Install above).

## Provision the application once

An administrator creates an application using the monitoring API. This helper is
included in the same package:

```python
import os
from aictrl_langgraph import register_application

credentials = register_application(
    base_url=os.environ["AICTRL_URL"],
    name="Support workflow",
    management_token=os.environ["MONITORING_ACCESS_TOKEN"],
)
# Store credentials.api_key in your secret manager as AICTRL_API_KEY.
# Save credentials.application_id for administration. Do not log the key.
```

Omit `management_token` only for the backend's explicit local development mode.
Do not provision on every process start. Registration is not retried because a
lost response can already have created an application. If this happens, inspect
the application list and rotate that application's credentials through the API.

Runtime processes only need `AICTRL_URL` and the application-scoped `AICTRL_API_KEY`.
`AICTRL_URL` is the server URL (for example `http://localhost:8100` locally), optionally
ending in `/api/v1`. Use HTTPS for remote connections. The SDK does not load `.env`
files. A management token is not an ingestion key.

## Integrate your graph

```python
from aictrl_langgraph import LangGraphMonitor

# builder is your application's existing StateGraph.
with LangGraphMonitor.from_env() as monitor:
    graph = monitor.instrument(builder.compile())
    result = graph.invoke(inputs)
    monitor.flush()
```

`instrument` returns a copy of the compiled graph with monitoring callbacks added.
Existing callbacks, configuration, checkpointing, runtime context and graph methods
remain available. Use the returned graph for `invoke`, `batch`, `stream`, `ainvoke`,
`abatch` and `astream`. Use one monitor per root graph and application connection;
create a new monitor for a new graph version. Instrument during startup.

Async applications can keep network registration and shutdown off the event loop:

```python
async with LangGraphMonitor.from_env() as monitor:
    graph = await monitor.ainstrument(builder.compile())
    result = await graph.ainvoke(inputs)
    async for chunk in graph.astream(inputs, stream_mode="updates"):
        consume(chunk)
    await monitor.aflush()
```

Stream modes and chunks are unchanged. Exhaust or explicitly close/aclose stream
iterators before closing the monitor. Closing a stream early or cancelling an
invocation records a failed run. Graph exceptions are preserved.

A runnable, LLM-free example is included in
[`examples/workflow.py`](examples/workflow.py). Run it only against an application
registered for testing; it creates real monitoring records.

## Run and graph semantics

- Each graph invocation gets a new monitoring run ID and starts at sequence 1.
  Parallel invocations are isolated. Parallel nodes and retry attempts share the
  invocation's ordered sequence.
- Interrupts finish that invocation as `PAUSED`. Resuming a checkpoint creates a
  **new monitoring run** for the resumed invocation. LangGraph's thread/checkpoint
  continuity remains intact; thread IDs and checkpoint contents are not sent.
- Branches and loops retain their node IDs. Give conditional edges an explicit
  `path_map` or a `Literal` return annotation so LangGraph can expose their topology.
  The SDK exports the public `get_graph(xray=False)` representation; it cannot
  discover undeclared runtime `Command` destinations. Declare node destinations
  when building such graphs.
- Nested subgraphs appear as their parent node. Internal subgraph nodes, nested
  runnable calls, and individual model/tool calls are not separate runs or nodes.
  Cache hits that do not execute a node do not generate synthetic node events.
- Node IDs must match `[\w.:-]{1,100}`. Graphs allow at most 500 nodes and 2,000 edges.
  The immutable digest returned by the API is attached to every event.
- Monitoring callbacks observe execution without changing routing, prompts, retry
  policies or checkpoint state. Explicit policy guards can block node execution.

## Guard nodes with policies

Declare local evaluators with `Policy` and wrap each protected node with
`monitor.guard_node` **before** adding it to the StateGraph:

```python
from aictrl_langgraph import LangGraphMonitor, Policy, PolicyBlocked

policies = [
    Policy(
        "export_limit",
        "Export amount limit",
        "1",
        lambda state: state["amount"] <= 100,
        description="Allow exports up to 100 units.",
        operations=("data.export",),
    )
]
with LangGraphMonitor.from_env(policies=policies) as monitor:
    # builder, export_result and inputs belong to your application.
    builder.add_node(
        "export",
        monitor.guard_node(
            "export",
            export_result,
            operation="data.export",
        ),
    )
    graph = monitor.instrument(builder.compile())
    try:
        result = graph.invoke(inputs)
    except PolicyBlocked:
        handle_denial()
```

The catalog appears in **Policies**. Choose **Optimize**, assign a policy to a
diamond and **Apply**. Each guard fetches current assignments before calling the
node, so a saved change affects its next execution boundary. No assignment means
no evaluator runs, but the configuration read must still succeed. Bump the policy
version when its behavior changes; evaluator source code is never uploaded.

Exactly `True` allows execution; `ReviewRequired` pauses for human approval.
Denials and evaluator errors raise `PolicyBlocked` before the node function runs.
Unavailable configuration, incompatible active versions and review-service errors
raise its subclass `PolicyUnavailable`, allowing applications to distinguish a
connection/configuration problem from a policy refusal. Callbacks report both as
`BLOCKED`. Keep evaluators
free of side effects and do not mutate their input. Async node functions support
async evaluators; sync nodes reject an async evaluator. Wrappers preserve node
signatures and injected config/runtime context.

Guard only with ordinary sync/async node functions. Unwrapped nodes and side
effects already in progress are outside enforcement; parallel operations are not
rolled back. Do not cache protected nodes: cache hits skip their guards. A
checkpoint resume checks nodes that execute again. All versions using one
application key must have compatible graph/catalog metadata; incompatible active
bindings block until a compatible configuration is applied.

## Delivery and shutdown

```python
monitor = LangGraphMonitor(
    base_url="https://control.example.com",
    api_key=application_key,
    timeout=5.0,
    max_retries=2,
    queue_capacity=10_000,
    heartbeat_interval=10.0,
    max_active_runs=1_000,
)
```

Callbacks only enqueue metadata under a short lock. One background worker sends
FIFO batches of at most 100 events. Transient network errors and HTTP
408/429/500/502/503/504 are retried with the **same event IDs, timestamps and body**.
Other errors, including authentication failures and sequence conflicts, are not
retried. Redirects are not followed. Request timeouts and retry counts are bounded.

If delivery retries are exhausted or the queue fills, telemetry for that monitor
stops and `monitor.error` records the failure. It does not silently skip an event
and continue with an invalid sequence. The application keeps running, a safe
warning is logged, and `flush()` / `close()` raise `MonitoringError`. Opt-in guards
still require a successful configuration read before execution; an unavailable
connection denies the operation even if telemetry can no longer report it.
Telemetry failure leaves the HTTP client available for fresh policy reads and
human-review decisions once the service returns. Only explicit monitor closure
closes that client. Create a new monitor to restore telemetry after the underlying
problem is corrected. The affected history can remain incomplete; completion is
never fabricated.

`flush(timeout=30)` waits for queued events to be acknowledged. A flush timeout
does not discard events. `close(timeout=30)` stops new telemetry and drains the
queue; use it after all graph invocations and streams finish. `aflush` and `aclose`
are the asynchronous equivalents. Context managers close automatically and do
not replace an exception already raised by the application with a delivery error.

Delivery is in-memory and best effort. Process crashes or forced termination can
lose unsent events. There is no disk spool, cross-process queue or recovery of
in-flight invocations after a process restart. Create monitors after worker fork;
monitor instances belong to one process.

## Data collected

Only graph node IDs/labels, topology, opaque monitoring run/event IDs, graph digest,
lifecycle types, sequence numbers and UTC timestamps are sent. With policy guards,
policy IDs/names/descriptions/versions, operation restrictions and hook metadata are
also sent, and the SDK reads assignment IDs. Inputs, outputs,
state values, messages, prompts, tool arguments, exception text, checkpoint data
and callback metadata are never serialized. Node names and edge labels are visible
in the UI, so use structural labels rather than user content. Credentials and
server response bodies are excluded from SDK diagnostics.

## Development

```sh
uv run pytest packages/aictrl-langgraph/tests services/backend/tests/test_langgraph_sdk.py
make verify
make sdk-check
```

`make sdk-check` builds the wheel and source distribution, checks their metadata,
installs the wheel into an isolated environment and executes a real LangGraph run
outside the repository. CI performs isolated installation checks on Python
3.11, 3.12, 3.13 and 3.14. These checks do not upload or publish anything.

## Human-in-the-loop policies

Return `ReviewRequired` from a policy evaluator to require an editor's decision:

```python
from aictrl_langgraph import Policy, ReviewRequired

review = Policy(
    id="pii_review",
    name="Approve personal data release",
    version="1",
    operations=("chat.respond",),
    check=lambda state: ReviewRequired(state["request_id"]) if state["sensitive"] else True,
)
```

Declare this policy on the monitor, guard the response node, compile with a
LangGraph checkpointer, and instrument as usual. Assign it in the control-plane
policy editor. `request_id` must be an opaque unique identifier retained in the
checkpoint; never put user content in it. The guarded input must be JSON-compatible.
The SDK binds the review to the entire input using a keyed HMAC, scoped to the
application, node, policy and active graph/catalog/binding versions.

On pause, LangGraph's `__interrupt__` result contains a safe payload:
`{"type": "action_required", "review_ids": [...]}`. Forward that status to the
user while keeping the response on the backend. Poll `monitor.review_status(id)`
or `await monitor.areview_status(id)`: results are `pending`, `approved`, `rejected`
or `expired`. When all required reviews are approved, resume the same checkpoint
with `Command(resume=True)`. The guard fetches the service decision again and never
trusts the resume value. Reject/expiry/read failures block; modified input cannot
reuse the request identity. A block from any assigned evaluator wins over review.

The **Reviews** panel uses policy-editor access for decisions. Approvals expire
after ten minutes, and graph/catalog/binding changes invalidate them. Request and
policy metadata, keyed input digests and timestamps are stored separately from
lifecycle events; protected input never leaves the SDK through this API. Resume
still creates a new monitoring run. The application owns checkpoint retention,
request-to-conversation correlation, delivery idempotency and disconnect handling.
See the [SQLite chat example](../../apps/sqlite_chat/README.md) for a complete flow.
