Streaming LLMAgent Executor — Progress Updates During Task Execution¶
This notebook contrasts StreamingLLMAgentA2AExecutor with the
non-streaming LLMAgentA2AExecutor from examples/ch10.ipynb's
Example 4: the same Hailstone task, dispatched to two real, separately
running A2A servers, one publishing an incremental
TaskStatusUpdateEvent per step and one publishing a single terminal
update.
NOTE
NOTE
There is no client-side streaming support in UseA2AAgentTool yet;
dispatch there is always ClientConfig(streaming=False), one
ToolCallResult per call, per its own docstring. So this notebook
drives the raw a2a-sdk client directly.
Setup Instructions¶
This notebook launches the extra/a2a-from-scratch-hailstone and
extra/a2a-from-scratch-hailstone-streaming apps from a local checkout,
so it needs the full repository, not just the package. Launch it with:
uv run --with jupyter jupyter lab
from the project root. No PyPI release of llm-agents-from-scratch has
A2A support yet (still under "Unreleased" in CHANGELOG.md), so the
pip install cell below won't provide what this notebook needs even
if uncommented; it's left only for parity with the other notebooks.
Setting the backbone LLM of your agent¶
These notebooks run on Ollama by default, the setup the book teaches.
If you do nothing, nothing changes: make_llm() starts a local Ollama
service when one isn't already running.
To use OpenAI or Anthropic instead:
- Install the extra:
uv sync --extra openaioruv sync --extra anthropic. - Export
OPENAI_API_KEYorANTHROPIC_API_KEYbefore launching Jupyter. - Set
LLM_PROVIDER=openai(oranthropic), or passprovider="openai"tomake_llm(). A key on its own never switches providers, so one exported for unrelated work cannot reroute you off the Ollama path.
The switch applies wherever a notebook builds its LLM with make_llm().
A notebook that constructs OllamaLLM directly stays on Ollama regardless
of these settings.
If you opted in but forgot to export the key, you will be prompted for it.
That is the safe path on hosted kernels, and it keeps the key out of the
saved notebook. Setting OLLAMA_API_KEY alone routes Ollama to Ollama
Cloud.
Caveat: the examples are tuned for qwen3. Output on gpt-5 or Claude will differ from what is printed in the book, and prompt-sensitive examples (the ch09 evaluator pattern, ch08 supervised trajectories) may behave noticeably differently.
# No PyPI release has A2A support yet -- see Setup Instructions above.
# !pip install llm-agents-from-scratch
Running an Ollama service¶
To execute the code provided in this notebook, you'll need to have
Ollama installed on your local machine and have its LLM hosting
service running. To download Ollama, follow the instructions found on
this page: https://ollama.com/download. After downloading and
installing Ollama, you can start a service by opening a terminal and
running ollama serve.
import os
import subprocess
import time
import urllib.error
import urllib.request
from llm_agents_from_scratch.notebook_utils import ensure_ollama
use_cloud = "OLLAMA_API_KEY" in os.environ
ensure_ollama() if not use_cloud else print("\u2713 Using Ollama Cloud")
✓ Using Ollama Cloud
The Two Servers¶
Both peers are the same standalone apps examples/ch10.ipynb uses,
extra/a2a-from-scratch-hailstone and its sibling
extra/a2a-from-scratch-hailstone-streaming, genuinely separate OS
processes, not co-routines sharing this notebook's process. The only
difference between the two apps is which executor wraps the agent:
LLMAgentA2AExecutor versus StreamingLLMAgentA2AExecutor.
import contextlib
import signal
from pathlib import Path
from urllib.parse import urlparse
def find_repo_root(marker="extra"):
"""Walk up from cwd to the first ancestor containing `marker`.
Jupyter sets a kernel's cwd to the directory of the opened notebook
file under normal use, but that's not guaranteed for every launch
method (e.g. running via nbconvert from the project root), so this
doesn't assume a fixed number of parent hops.
"""
for candidate in (Path.cwd(), *Path.cwd().parents):
if (candidate / marker).is_dir():
return candidate
msg = f"Could not find a parent directory containing {marker!r}."
raise RuntimeError(msg)
def ensure_hailstone_server(app_dir, host, log_path, timeout=15):
"""Start a from-scratch Hailstone A2A server app if not already running.
Returns the Popen handle if this call started the server, or None if
it was already running (so Cleanup knows not to tear it down). Runs
in its own session (start_new_session=True) so Cleanup can kill the
whole process group -- terminating only the `uv run` wrapper process
can leave the actual uvicorn process orphaned and still bound to the
port.
"""
def _up():
try:
with urllib.request.urlopen(
f"{host}/.well-known/agent-card.json",
timeout=1,
):
pass
return True
except (urllib.error.URLError, ConnectionError, TimeoutError):
return False
if _up():
print(f"\u2713 {app_dir} A2A server already running at {host}")
return None
server_path = find_repo_root() / "extra" / app_dir
port = urlparse(host).port or 80
print(f"Starting {app_dir} A2A server at {host}...")
with open(log_path, "w") as log_file:
process = subprocess.Popen(
[
"uv",
"run",
"uvicorn",
"main:app",
"--host",
"0.0.0.0",
"--port",
str(port),
],
cwd=server_path,
stdout=log_file,
stderr=subprocess.STDOUT,
start_new_session=True,
)
deadline = time.time() + timeout
while time.time() < deadline:
if _up():
print(f"\u2713 {app_dir} A2A server up at {host}")
return process
time.sleep(0.5)
stop_hailstone_server(process)
raise RuntimeError(f"A2A server did not start within {timeout}s")
def stop_hailstone_server(process):
"""Kills the whole process group started by ensure_hailstone_server()."""
def _signal(sig):
with contextlib.suppress(ProcessLookupError):
os.killpg(os.getpgid(process.pid), sig)
_signal(signal.SIGTERM)
try:
process.wait(timeout=5)
except subprocess.TimeoutExpired:
_signal(signal.SIGKILL)
process.wait(timeout=5)
def cleanup_servers():
"""Stops both servers this notebook started, ignoring already-up ones.
Looked up by name at call time rather than passed in, so later
cells (discovery, dispatch) can call this from an except block on
the same two module-level names Cleanup uses, without threading
the process handles through every intervening cell.
"""
for process in (streaming_process, nonstreaming_process):
if process is not None:
stop_hailstone_server(process)
started_servers = []
try:
streaming_process = ensure_hailstone_server(
"a2a-from-scratch-hailstone-streaming",
"http://127.0.0.1:9301",
"streaming_hailstone_server.log",
)
started_servers.append(streaming_process)
nonstreaming_process = ensure_hailstone_server(
"a2a-from-scratch-hailstone",
"http://127.0.0.1:9300",
"nonstreaming_hailstone_server.log",
)
started_servers.append(nonstreaming_process)
except Exception:
for started in started_servers:
if started is not None:
stop_hailstone_server(started)
raise
Starting a2a-from-scratch-hailstone-streaming A2A server at http://127.0.0.1:9301... ✓ a2a-from-scratch-hailstone-streaming A2A server up at http://127.0.0.1:9301 Starting a2a-from-scratch-hailstone A2A server at http://127.0.0.1:9300... ✓ a2a-from-scratch-hailstone A2A server up at http://127.0.0.1:9300
Dispatching with a Raw Streaming Client¶
A2AAgentSpec.from_url() still does discovery, the same as
examples/ch10.ipynb. Everything after that is hand-rolled: build an
SDK client with ClientConfig(streaming=...), send one message, and
print every chunk client.send_message() yields as it arrives, rather
than collapsing to a single final result the way UseA2AAgentTool
does.
import httpx
from a2a.client import ClientConfig, create_client
from a2a.helpers import new_text_message
from a2a.types import Role as A2ARole
from a2a.types import SendMessageRequest, TaskState
from llm_agents_from_scratch.a2a import A2AAgentSpec
from llm_agents_from_scratch.a2a.client.utils import a2a_parts_text
try:
streaming_spec = await A2AAgentSpec.from_url(
"http://127.0.0.1:9301",
timeout=120.0,
)
nonstreaming_spec = await A2AAgentSpec.from_url(
"http://127.0.0.1:9300",
timeout=120.0,
)
except Exception:
cleanup_servers()
raise
async def dispatch_and_print(spec, streaming, task_text):
"""Send one task to a peer, printing each StreamResponse chunk live."""
async with httpx.AsyncClient(timeout=120.0) as httpx_client:
client = await create_client(
agent=spec.agent_card,
client_config=ClientConfig(
streaming=streaming,
httpx_client=httpx_client,
),
)
n_chunks = 0
try:
message = new_text_message(text=task_text, role=A2ARole.ROLE_USER)
request = SendMessageRequest(message=message)
async for chunk in client.send_message(request):
n_chunks += 1
kind = chunk.WhichOneof("payload")
if kind == "status_update":
state = TaskState.Name(chunk.status_update.status.state)
has_message = chunk.status_update.status.HasField("message")
text = (
a2a_parts_text(chunk.status_update.status.message.parts)
if has_message
else ""
)
print(f" [{n_chunks}] {state}: {text[:90]!r}")
elif kind == "task":
state = TaskState.Name(chunk.task.status.state)
print(f" [{n_chunks}] task created, state={state}")
elif kind == "artifact_update":
text = a2a_parts_text(chunk.artifact_update.artifact.parts)
print(f" [{n_chunks}] artifact_update: {text[:90]!r}")
finally:
await client.close()
print(f" total chunks: {n_chunks}")
return n_chunks
Example — The Same Task, Two Executors¶
6 is a short enough Hailstone sequence to run in seconds, but long enough (8 steps) to show several distinct status updates land before completion.
task_text = (
"Compute the hailstone sequence starting at 6, step by step using "
"next_number, until it reaches 1."
)
try:
print("=== Streaming ===")
streaming_chunks = await dispatch_and_print(streaming_spec, True, task_text)
print()
print("=== Non-streaming ===")
nonstreaming_chunks = await dispatch_and_print(
nonstreaming_spec,
False,
task_text,
)
except Exception:
cleanup_servers()
raise
=== Streaming ===
[1] task created, state=TASK_STATE_SUBMITTED
[2] TASK_STATE_SUBMITTED: ''
[3] TASK_STATE_WORKING: ''
[4] TASK_STATE_WORKING: 'Compute the hailstone sequence starting at 6, step by step using next_number, until it rea'
[5] TASK_STATE_WORKING: 'I need to make the following tool-calls:\n{\n "id_": "ba0e88df-b8f3-4c70-bfbe-dce741d00cb'
[6] TASK_STATE_WORKING: 'Make the next tool call: next_number with x=3. The previous call returned 3, so continue c'
[7] TASK_STATE_WORKING: 'I need to make the following tool-calls:\n{\n "id_": "e6544ee9-e84a-4034-834f-756108e522b'
[8] TASK_STATE_WORKING: 'The current_response says it needs to make the tool call next_number with x=10. Execute th'
[9] TASK_STATE_WORKING: 'I need to make the following tool-calls:\n{\n "id_": "2963ce52-0579-4d29-9578-dab0435d714'
[10] TASK_STATE_WORKING: 'Make the next tool call: next_number with x=5. Continue computing the hailstone sequence s'
[11] TASK_STATE_WORKING: 'I need to make the following tool-calls:\n{\n "id_": "585e8aea-c04a-4628-9d69-6d47035e7f9'
[12] TASK_STATE_WORKING: 'Make the next tool call: next_number with x=16. Continue computing the hailstone sequence '
[13] TASK_STATE_WORKING: 'I need to make the following tool-calls:\n{\n "id_": "7be1e405-3e8f-49f1-8894-4d4bd3bb248'
[14] TASK_STATE_WORKING: 'Make the next tool call: next_number with x=8. Continue computing the hailstone sequence s'
[15] TASK_STATE_WORKING: 'I need to make the following tool-calls:\n{\n "id_": "182f97ca-bf2c-41a9-8192-7070122cbcc'
[16] TASK_STATE_WORKING: 'Make the next tool call: next_number with x=4. Continue computing the hailstone sequence s'
[17] TASK_STATE_WORKING: 'I need to make the following tool-calls:\n{\n "id_": "cda1ddc7-90f7-44bf-bcce-a8760c83ce9'
[18] TASK_STATE_WORKING: 'Make the next tool call: next_number with x=2. Continue computing the hailstone sequence s'
[19] TASK_STATE_WORKING: 'The hailstone sequence has reached 1. The complete sequence starting from 6 is:\n\n**6 → 3 →'
[20] artifact_update: 'The hailstone sequence has reached 1. The complete sequence starting from 6 is:\n\n**6 → 3 →'
[21] TASK_STATE_COMPLETED: ''
total chunks: 21
=== Non-streaming ===
[1] task created, state=TASK_STATE_COMPLETED
total chunks: 1
The streaming peer's chunk count includes two updates per step
(the planned instruction, then the step's result), plus the initial
task-submitted chunk, a trailing artifact, and the final completed
status; a modest 8-step sequence already produces around 20 chunks.
The non-streaming peer produces exactly one, the same shape
UseA2AAgentTool already gets today.
Cleanup¶
cleanup_servers()