Metadata-Version: 2.4
Name: ag-ui-crewai
Version: 0.2.2a3
Summary: Implementation of the AG-UI protocol for CrewAI
Author-email: Markus Ecker <markus.ecker@gmail.com>
License-Expression: MIT
License-File: LICENSE
Requires-Python: <3.14,>=3.10
Requires-Dist: ag-ui-a2ui-toolkit>=0.0.4
Requires-Dist: ag-ui-protocol>=0.1.19
Requires-Dist: crewai<2,>=1.0
Requires-Dist: fastapi<0.116.0,>=0.115.12
Requires-Dist: litellm>=1.60.2
Requires-Dist: uvicorn<0.35.0,>=0.34.3
Provides-Extra: tools
Requires-Dist: crewai-tools>=0.38.1; extra == 'tools'
Description-Content-Type: text/markdown

# ag-ui-crewai

Implementation of the AG-UI protocol for CrewAI.

Provides a complete Python integration for CrewAI flows and crews with the AG-UI protocol, including FastAPI endpoint creation and comprehensive event streaming.

## Installation

```bash
pip install ag-ui-crewai
```

## Usage

```python
from crewai.flow.flow import Flow, start
from litellm import acompletion
from ag_ui_crewai import (
    add_crewai_flow_fastapi_endpoint,
    copilotkit_stream,
    CopilotKitState
)
from fastapi import FastAPI

class MyFlow(Flow[CopilotKitState]):
    @start()
    async def chat(self):
        response = await copilotkit_stream(
            await acompletion(
                model="openai/gpt-4o",
                messages=[
                    {"role": "system", "content": "You are a helpful assistant."},
                    *self.state.messages
                ],
                tools=self.state.copilotkit.actions,
                stream=True
            )
        )
        self.state.messages.append(response.choices[0].message)

# Add to FastAPI
app = FastAPI()
add_crewai_flow_fastapi_endpoint(app, MyFlow(), "/flow")
```

### Conversational Flows

CrewAI 1.15.11's Conversational Flows use the same AG-UI event translation,
state synchronization, tools, reasoning, multimodal content, interrupts, and
generative UI support as regular Flows. Opt the Flow into CrewAI's public
conversation API and register the endpoint with `conversational=True`.

> **Important:** CrewAI builds a Flow's graph from the subclass's own
> `__dict__`, so a subclass that only sets `conversational = True` inherits
> **none** of the base Flow's `@start`/`@listen` methods and runs an empty graph
> (your steps silently never fire). Re-copy the base's flow methods onto the
> conversational type — this is exactly what the dojo's
> `examples/conversational.py::_conversational_type` helper does:

```python
from crewai.experimental.conversational import ConversationConfig

_flow_methods = {
    name: value
    for name, value in MyFlow.__dict__.items()
    if not name.startswith("_") and hasattr(value, "__flow_method_definition__")
}

MyConversationalFlow = type(
    "MyConversationalFlow",
    (MyFlow,),
    {
        **_flow_methods,
        "conversational": True,
        "conversational_config": ConversationConfig(defer_trace_finalization=False),
    },
)

add_crewai_flow_fastapi_endpoint(
    app,
    MyConversationalFlow(),
    "/conversational-flow",
    conversational=True,
)
```

The bridge invokes `flow.stream_turn(message, session_id=thread_id)`: AG-UI's
`threadId` is the CrewAI conversation session ID. It hydrates prior messages into
the Flow state before the current turn, passes only the latest user's text to
`stream_turn`, and preserves media blocks on that current message. Each HTTP
request finalizes its own CrewAI trace even if the Flow's conversation config
would normally defer finalization across turns.

Conversational mode requires CrewAI's ordered `StreamFrame` transport and a Flow
that both sets `conversational=True` and exposes `stream_turn`. If any part of
that contract is unavailable, the endpoint emits a correlated `RUN_ERROR` with
code `AGUI_CREWAI_CONVERSATIONAL_FLOW_UNSUPPORTED`; it never silently falls back
to a regular Flow kickoff. `get_capabilities()["conversationalFlows"]` declares
whether the installed runtime exposes both the required transport and public
turn API.

The AG-UI dojo presents these as two separate framework choices:

- `crewai`: **CrewAI Flows**, preserving the existing `/crewai/...` URLs and
  including the legacy `crew_chat` example.
- `crewai-conversational-flows`: **CrewAI Conversational Flows**, with the same
  Flow feature matrix under `/crewai-conversational-flows/...`; `crew_chat` is
  intentionally excluded because it is not a Flow.

## Features

- **Native CrewAI integration** – Direct support for CrewAI flows, crews, and multi-agent systems
- **FastAPI endpoint creation** – Automatic HTTP endpoint generation with proper event streaming
- **Predictive state updates** – Real-time state synchronization between backend and frontend
- **Streaming tool calls** – Live streaming of LLM responses and tool execution to the UI
- **Backend tool rendering** – Tools bound to a CrewAI `Agent`/`Crew` run server-side and surface to the UI as a tool call plus a `TOOL_CALL_RESULT`, so the client can render them without executing the tool (see the `backend_tool_rendering` example). Requires the StreamFrame transport (crewai >= 1.6); on crewai 1.0–1.5 the legacy event-bus path does not surface backend tool calls. A tool that returns structured data should return it as a JSON string (e.g. `json.dumps(...)`), since crewai stringifies tool output before it reaches the bridge.

## Protocol surface

### Wire shape: START / CONTENT / END triples (default)

Text and tool-call output is emitted as `TEXT_MESSAGE_START` / `TEXT_MESSAGE_CONTENT`
/ `TEXT_MESSAGE_END` and `TOOL_CALL_START` / `TOOL_CALL_ARGS` / `TOOL_CALL_END`, the
protocol's canonical discrete form. `emission_shape="chunks"` (or
`AGUI_CREWAI_EMISSION_SHAPE=chunks`) opts back into the previous
`TEXT_MESSAGE_CHUNK` / `TOOL_CALL_CHUNK` form.

```python
add_crewai_flow_fastapi_endpoint(app, MyFlow(), "/flow")                     # triples
add_crewai_flow_fastapi_endpoint(app, MyFlow(), "/flow", emission_shape="chunks")
```

The two shapes are **not** equivalent. In `chunks` mode a `copilotkit_emit_state` /
`copilotkit_predict_state` call that lands between two tool-call argument deltas makes
`@ag-ui/client`'s chunk transform close the call and then throw on the next delta;
triples keep the call open server-side, because the server knows more deltas are
coming and the client does not. Triples are also the form any consumer can apply
directly (`apply/default.ts` throws if a chunk reaches it untransformed), so a raw
SSE reader (the Python SDK, conformance tooling, custom clients) needs no
chunk-transform stage.

Both transports (the crewai >= 1.6 `StreamFrame` path and the legacy
event-bus-listener fallback) route through one `EmissionShaper`, so the event shape
and payload never depend on the installed crewai version. A run never ends with an
open sequence: any open message, tool call, or step is closed before `RUN_FINISHED`.
MCP tool executions always use triples regardless of this setting: their name, args
and result arrive together rather than streamed.

### RAW passthrough (opt-in, default OFF)

`emit_raw_events=True` mirrors the crewai events this bridge does **not** map onto
AG-UI `RAW` events: crewai's llm / agent / task / tool channels (including
`llm_thinking_chunk`), nested-flow lifecycle events, and internals such as `cc_env`.
Events the bridge already maps are never duplicated as RAW.

```python
add_crewai_flow_fastapi_endpoint(app, MyFlow(), "/flow", emit_raw_events=True)
# or, without a code change:  AGUI_CREWAI_EMIT_RAW_EVENTS=1
```

It is off by default deliberately: LangGraph shipped RAW passthrough on and the
payload bloat had to be walked back. RAW payloads are large and can carry prompt and
completion text, so enabling it widens what leaves your process.

Requires the `StreamFrame` transport: the installed crewai must expose it (>= 1.6)
**and** the served flow must expose `astream`, which the driver probes per flow. On
the legacy event-bus fallback the bridge logs one warning per process and emits no
RAW events, because that listener never sees the unmapped events.

RAW mirrors never precede `RUN_STARTED`. crewai raises some events before the flow
opens, and a `RAW` first event makes the reference client reject the whole stream, so
those are held and released once the run has opened. Both RAW buffers are bounded; a
saturated buffer degrades mirroring (logged) rather than the run.

### Memory is isolated per `threadId` (default ON)

A crew served with `Crew(memory=True)` keeps its memories in one on-disk store,
namespaced by the _crew name_. Nothing in that namespace derives from the AG-UI
`threadId`, so without help every chat served by an endpoint reads and writes the
same namespace and one user's remembered facts surface in another user's chat.
(Setting `inputs["id"] = thread_id` does not help: that scopes crewai's flow-state
persistence, a different subsystem.) `Agent(memory=True)` has the same shape one
level down: the agent builds its _own_ memory, which crewai prefers over the
crew's.

The bridge closes that by giving each request a `MemoryScope` view of the crew's
memory — and of each agent's own memory — rooted at a path derived from the
request's `threadId`. Threads are mutually invisible; each still sees its own
history across sequential runs. One physical store, no directory-per-thread
sprawl.

Because crewai picks the executing agent off `task.agent` (or `manager_agent`
under the hierarchical process) and reaches the crew's memory through
`agent.crew`, the request gets shallow _views_ of the crew, its agents and its
tasks, wired to each other. Nothing shared between concurrent requests is
mutated, and everything below the views (tools, LLMs, knowledge, the store
itself) stays shared.

```bash
# Opt out: restore the pre-fix behaviour of one memory namespace per crew,
# shared by every chat. Useful when the crew is a durable knowledge base
# rather than a per-conversation memory.
AGUI_CREWAI_THREAD_SCOPED_MEMORY=false
```

Limitations, in order of how likely you are to hit them:

- **Only crews and agents the bridge can reach are scoped.** That means the crew
  you passed to `add_crewai_crew_fastapi_endpoint`, plus any crew or standalone
  agent your `Flow` holds as an attribute (a crew's own agents and tasks come
  with it). A crew or agent _constructed inside_ a flow method is created after
  this point and is not scoped; construct it as a flow attribute, or pass it a
  `Memory` you scope yourself.
- **Per-request views are shallow.** Each request runs against copies of the
  crew, its agents and its tasks, so `crew.tasks[0].output` on the object you
  built is not filled in by a bridge-served run; read the run's result off the
  AG-UI event stream instead.
- **Isolation is logical, not physical.** All threads share one store and one
  embedder; a scope keeps reads and writes inside a namespace, it is not a
  security boundary against code that queries the store directly.
- **Older crewai degrades rather than crashing.** The bridge probes for crewai's
  unified memory view API at runtime. On a build without it, isolation is not
  active and the bridge logs one warning saying exactly that.

### `get_capabilities()`

Returns a capability declaration (the CrewAI counterpart of
`LangGraphAgent.get_capabilities`). `crewai.__version__` appears only as
informational metadata, never as a gate.

```python
from ag_ui_crewai import get_capabilities

get_capabilities(llm=my_agent.llm, emit_raw_events=True)
```

`transport`, `rawEvents`, `reasoning`, `conversationalFlows`, and `crewChat` come
from runtime probes; `humanInTheLoop` and `state` are static declarations of what
the bridge implements today. `emit_raw_events` defaults to re-reading the
environment, so pass the same value your endpoint was registered with if you want
the declaration to describe _that_ endpoint.

Reasoning surfaces as first-class `REASONING_*` events (`REASONING_START` /
`REASONING_MESSAGE_START` / `REASONING_MESSAGE_CONTENT` / `REASONING_MESSAGE_END` /
`REASONING_END`, plus `REASONING_ENCRYPTED_VALUE` for signature / redacted-thinking
blocks), **provider-agnostic** and on **both** transports. It needs neither
`emit_raw_events` nor the `StreamFrame` transport. Two channels feed it:

- **litellm delta** (`copilotkit_stream`): reads `reasoning_content` /
  `thinking_blocks` for any reasoning-capable model routed through litellm
  (deepseek-reasoner, Anthropic extended thinking, Bedrock, xAI,
  gemini-via-litellm, and reasoning models normalised by litellm). This is the
  provider-agnostic path and drives the crew-serving path in `crews.py` too.
- **native `LLMThinkingChunkEvent`** (crewai's Gemini provider, crewai >= 1.10.1):
  an additional source on the `StreamFrame` path.

`reasoning.supported` is therefore True whenever a reasoning channel is live (the
litellm channel is effectively always live, as litellm is a direct dependency). A
non-reasoning model simply emits nothing (graceful no-op). `requiresEmitRawEvents`
is `False`. The `nativeGeminiProvider` / `resolvedProvider` fields are
informational (the native event is an extra source, not a requirement), and
`thinkingEventAvailable` reports whether the native Gemini event resolved.

## Tuning knobs

The CrewAI integration exposes three environment variables for tuning
timeouts and teardown behaviour. Sensible defaults ship with the
package; override these only if your deployment has specific needs
(long-running crews, disconnect-heavy workloads, flaky LLM providers).

### `AGUI_CREWAI_LLM_TIMEOUT_SECONDS`

Per-read timeout forwarded to `litellm.acompletion` in
`ChatWithCrewFlow.chat`. It applies to all three completion sites: the
initial call, the post-crew-run follow-up (tool-choice=`"none"`) that
lets the assistant speak about the crew result, and the post-`crew_exit`
(tool-choice=`"none"`) call.

> **Limitation — no tool chaining after a crew run.** The post-crew-run
> follow-up uses `tool_choice="none"`, so the assistant summarizes the crew
> result as text but cannot call a frontend action in the same turn. A flow
> like "run the crew, then update the UI" is not reachable on this path
> today; allowing bounded tool re-entry there is future work.

- **Default:** `120` seconds.
- **Non-positive** (e.g. `0`, `-1`): disables the per-read timeout —
  the underlying HTTP client's default applies instead.
- **Non-finite** (`nan`, `inf`): falls back to the default.
- **Note:** LiteLLM forwards this as a **per-read** timeout to the
  underlying HTTP client, not a session-level ceiling. A trickle-feeding
  server can keep the coroutine alive indefinitely at this layer; use
  `AGUI_CREWAI_FLOW_TIMEOUT_SECONDS` for the session-level cap.

### `AGUI_CREWAI_FLOW_TIMEOUT_SECONDS`

Hard wall-clock ceiling on a single flow run. Guards against a runaway
flow (hung LiteLLM stream, infinite loop in a user task) pinning the
process indefinitely.

- **Default:** `600` seconds (10 minutes).
- **Non-positive**: disables the ceiling. Only use this for
  deployments with legitimately long-running crews where the wall-clock
  ceiling is handled at a higher layer.
- **Non-finite** (`nan`, `inf`): falls back to the default.
- When the ceiling fires, the stream yields a `RUN_ERROR` event with
  code `AGUI_CREWAI_FLOW_TIMEOUT` and a message carrying the configured
  ceiling plus thread/run correlation IDs.
- **Conversational mode:** the ceiling bounds the AG-UI HTTP response, not
  the CrewAI worker. Conversational Flows drive a synchronous `StreamSession`
  on a background thread that cannot be closed from the request loop, so a
  hung upstream call keeps that worker alive until it emits or returns. (The
  async StreamFrame path cancels the CrewAI kickoff task and pins nothing; this
  opt-in conversational path does not.)

### `AGUI_CREWAI_CANCEL_JOIN_TIMEOUT_SECONDS`

Teardown ceiling: the total wall-clock budget for `_cancel_and_join` to
unwind the kickoff task after a client disconnect, timeout, or error.
Covers the grace window, force-cancel join, AND outer-cancel recovery
— one shared monotonic deadline, not three.

- **Default:** `10` seconds.
- **Non-positive** or **non-finite**: falls back to the default
  (deliberately not disable-able — a cancel that cannot be bounded is a
  resource leak).
- Tune upward if your deployment sees disconnect-heavy load and a
  consistently-stuck cancel warning is logged.

## To run the dojo examples

```bash
cd integrations/crew-ai/python
uv sync
uv run dev
```
