Background Execution
Run agents in the background. Reconnect to in-progress streams via SSE.
Run agents, teams, and workflows in the background by passing background=True to .arun(). Execution continues even if the client disconnects. The behavior depends on whether you also set stream=True.
Execution Modes
background | stream | Behavior |
|---|---|---|
False | True | Default streaming. Runs inline. Client disconnect cancels the run. |
False | False | Non-streaming. Returns full response. |
True | False | Fire-and-forget. Returns PENDING immediately. Poll for results. |
True | True | Resumable streaming. Runs in a detached task. Events are buffered. Reconnect via /resume. |
Background execution requires a database (db) on the agent, team, or workflow for persisting run state.
Setup
In an activated virtual environment:
uv pip install -U "agno[os]" openai "psycopg[binary]" sqlalchemy yfinance httpx
export OPENAI_API_KEY="your_openai_api_key"
mkdir -p tmpThe Agent and Team recipes require a reachable PostgreSQL database at the shown URL; replace it with your own connection URL. The Workflow recipe uses a writable local SQLite file. Save one recipe to a Python file and run it with python <filename>.py.
Fire-and-Forget
Start a background run and poll for the result. Works identically for agents, teams, and workflows.
import asyncio
from agno.agent import Agent
from agno.db.postgres import PostgresDb
from agno.models.openai import OpenAIResponses
from agno.run.base import RunStatus
db = PostgresDb(
db_url="postgresql+psycopg://ai:ai@localhost:5532/ai",
session_table="background_exec_sessions",
)
agent = Agent(
name="BackgroundAgent",
model=OpenAIResponses(id="gpt-5-mini"),
db=db,
)
async def main():
# Returns immediately with PENDING status
run_output = await agent.arun(
"Write a short analysis of quantum computing trends.",
background=True,
)
print(f"Run ID: {run_output.run_id}, Status: {run_output.status}")
# Poll until complete
for _ in range(60):
await asyncio.sleep(1)
result = await agent.aget_run_output(
run_id=run_output.run_id,
session_id=run_output.session_id,
)
if result and result.status == RunStatus.completed:
print(f"Done: {result.content}")
break
if result and result.status in (RunStatus.error, RunStatus.cancelled, RunStatus.paused):
print(f"Run stopped: {result.status}, {result.content}")
break
else:
print("Polling timed out. Exiting this script cancels its unfinished background task.")
asyncio.run(main())These SDK examples keep work in the current event loop. If polling times out and asyncio.run(main()) exits, loop shutdown cancels the unfinished task; saving IDs does not keep it executing. Keep the loop alive to continue the work. For client-independent execution, run a persistent AgentOS server; use its durable queue when work must survive server restarts.
Resumable Streaming (SSE)
Combine background=True with stream=True for resumable SSE streaming. The run executes in a detached asyncio.Task that survives client disconnects. Events are buffered with sequential event_index values so clients can reconnect and replay retained events.
Each AgentOS process retains the latest 10,000 events per run. Reconnect before the events you need are trimmed. A client that falls behind the retained window receives the remaining buffered events, but cannot recover the trimmed events from that in-memory buffer.
How It Works
Client connects → StreamingResponse reads from queue ← Background task runs
Client disconnects → StreamingResponse cancelled ← Background task keeps running
Client reconnects → /resume reads from subscriber queue ← Background task still publishing- The run persists
RUNNINGstatus in the database - A detached
asyncio.Taskexecutes and publishes events to an in-memory buffer - The client receives SSE events, each containing an
event_indexandrun_id - On disconnect, the client records
last_event_index - On reconnect, the client calls
/resumewithlast_event_indexto replay retained events after that index
Starting a Resumable Stream
Resumable streaming requires a running AgentOS server. Pass background=true and stream=true in the request. The pattern is the same for agents, teams, and workflows. Only the URL path differs.
Workflows also support WebSocket-based reconnection. See the WebSocket reconnect example.
import asyncio
import json
import httpx
BASE_URL = "http://localhost:7777"
async def start_resumable_stream():
async with httpx.AsyncClient(base_url=BASE_URL, timeout=60) as client:
# Use /agents, /teams, or /workflows
agents = (await client.get("/agents")).json()
agent_id = agents[0]["id"]
form_data = {
"message": "Write a detailed story about a brave knight.",
"stream": "true",
"background": "true",
}
run_id = None
session_id = None
last_event_index = None
async with client.stream("POST", f"/agents/{agent_id}/runs", data=form_data) as response:
buffer = ""
async for chunk in response.aiter_text():
buffer += chunk
while "\n\n" in buffer:
event_str, buffer = buffer.split("\n\n", 1)
for line in event_str.strip().split("\n"):
if not line.startswith("data: "):
continue
data = json.loads(line[6:])
# Track identifiers for reconnection
if data.get("run_id") and not run_id:
run_id = data["run_id"]
if data.get("session_id") and not session_id:
session_id = data["session_id"]
if data.get("event_index") is not None:
last_event_index = data["event_index"]
print(f"[{data.get('event_index')}] {data.get('event')}: {str(data.get('content', ''))[:60]}")
return run_id, session_id, last_event_index
asyncio.run(start_resumable_stream())Each SSE event includes:
event_index: Sequential integer for ordering and resumptionrun_id: The run identifier for reconnectionsession_id: The session identifier
Reconnecting via /resume
On disconnect (page refresh, network loss), reconnect to /resume with the last event_index:
async def resume_stream(agent_id: str, run_id: str, session_id: str, last_event_index: int):
form_data = {"last_event_index": str(last_event_index)}
if session_id:
form_data["session_id"] = session_id
async with httpx.AsyncClient(base_url=BASE_URL, timeout=120) as client:
async with client.stream(
"POST", f"/agents/{agent_id}/runs/{run_id}/resume", data=form_data
) as response:
buffer = ""
async for chunk in response.aiter_text():
buffer += chunk
while "\n\n" in buffer:
event_str, buffer = buffer.split("\n\n", 1)
for line in event_str.strip().split("\n"):
if not line.startswith("data: "):
continue
data = json.loads(line[6:])
event_type = data.get("event")
if event_type in ("catch_up", "replay", "subscribed"):
print(f"[META] {event_type}: {data}")
else:
print(f"[{data.get('event_index')}] {event_type}: {str(data.get('content', ''))[:60]}")Resume Endpoints
The resume endpoint follows the same pattern for agents, teams, and workflows:
POST /agents/{agent_id}/runs/{run_id}/resume
POST /teams/{team_id}/runs/{run_id}/resume
POST /workflows/{workflow_id}/runs/{run_id}/resume
Content-Type: application/x-www-form-urlencoded
last_event_index=N&session_id=SResume behavior depends on run state:
| Scenario | Condition | Behavior |
|---|---|---|
| Catch up + live | Run still active in this process's buffer | Replays retained events after last_event_index, then streams live events |
| Replay | Run completed, errored, cancelled, or paused and remains in this process's buffer | Replays retained events after last_event_index |
| DB fallback | Run absent from this process's buffer and session_id is provided | Replays persisted events after last_event_index when their original indices were stored; legacy unstamped events use positional replay |
Database replay requires store_events=True on the entity; exclusions still apply. Stored indices retain gaps and are filtered against last_event_index. Legacy events without stored indices receive positional indices and are not filtered by a live-stream index. Replaying persisted events does not restart execution.
If a run is absent from the buffer and session_id is missing, /resume returns an error. Completed, errored, and cancelled buffers become eligible for cleanup 30 minutes after finalization. An opportunistic cleanup check evicts eligible buffers when another run status is finalized, so eviction can occur later than 30 minutes.
Meta Events
The /resume stream can include these meta events:
| Event | Meaning |
|---|---|
catch_up | Run still active. Retained buffered events follow, then live events. |
replay | Run is no longer live in this process. Retained buffer events or persisted database events follow. |
subscribed | Catch-up is complete for an agent, team, or workflow. The client is now subscribed to live events. |
error | Run not found or other issue. |
Multi-Container Deployments
AgentOS can remove the affinity requirement entirely. QueueConfig(redis=...) wires a shared event stream and cancellation manager so /resume and /cancel work from any replica, and QueueConfig(durable=True) makes accepted runs survive process restarts. See AgentOS Background Execution.
The detached task and event buffer live in-process on the instance that started the run. The default cancellation manager is also in-memory. In a multi-replica setup, a /resume request that lands on a different instance can replay persisted events for a local entity when session_id is provided, then closes. It cannot tail the live task on the originating instance. A /cancel request on the wrong instance does not stop the task on the origin.
Configure session affinity on the initial run-start request, then route that client's /resume and /cancel requests to the same instance. An explicit run_id-to-instance mapping can provide the same guarantee. Hashing only the generated run_id on reconnect is insufficient because the initial request was routed before that ID existed. In AgentOS, QueueConfig(redis=...) removes the affinity requirement for both: it sets the cancellation manager (RedisRunCancellationManager) and the event stream (RedisEventStream) on every replica, so /cancel and live /resume work from any instance. See Multi-replica deployments.