Prerequisite Bootcamp for Agent Frameworks (LangGraph, Strands, ADK)
From core Python to typing, Pydantic, async, and tooling behind modern agent architectures.
uv, Pydantic v2, and framework docs.Core Focus Pillars
TypedDict, Annotated, Protocol, genericsBaseModel, TypeAdapterTaskGroup, cancellation, token streaminguv, ruff, pytest, strict CI quality gates| Day | Focus Area | Modules Included |
|---|---|---|
| Day 1 | Foundation | 1: Tooling with uv2: Core Language Fluency |
| Day 2 | Types | 3: Type Hints in Depth 4: Pydantic Validation & Schemas |
| Day 3 | Async | 5: Async and Concurrency 6: Framework Design Patterns |
| Day 4 | Real Work | 7: I/O, Data, and Logging 8: Testing and Quality |
| Day 5 | Ship | 9: Performance & Memory 10: Packaging, Security & Capstone |
Every module concludes with an executable hands-on lab and self-check criteria.
TypedDict, dataclasses, Protocol, and PEP 695 generics.uv, ruff, static type checking, pytest, and automated CI gates.if, for, while) and basic functions.clone, branch, commit, push).| Python Feature | LangGraph | Strands | ADK |
|---|---|---|---|
| TypedDict, Annotated | State channels, reducers | Tool parameters | State schemas |
| Pydantic models | Optional state validation | Structured output | Input/output schemas |
| Dataclasses | context_schema |
Config objects | Config & context |
| Decorators | @tool, @task, @entrypoint |
@tool, @hook |
@node |
| async, await, gather | ainvoke, astream |
stream_async |
run_node, gather |
| Async generators | Graph streaming | Streaming events | Event yields |
| Context managers | Checkpointer setup | MCPClient |
Runners, clients |
| Exceptions | GraphRecursionError |
Stop reasons, errors | BaseException trap |
Level 1 · Recognize
You can read the code and explain what it does in context.
Level 2 · Use
You can write the code from an example with documentation open.
Level 3 · Design
You choose the right abstraction and structure components cleanly.
Level 4 · Debug & Teach
You diagnose subtle runtime/concurrency bugs and mentor peers.
Target Benchmark: Reach Level 3 for typing, Pydantic, and async; reach Level 2 for tooling and patterns.
Python 3.10 – 3.11
requires-python target in legacy projects.TaskGroup.Python 3.12 – 3.13
asyncio.TaskGroup built-in.def f[T]()).itertools.batched, ReadOnly.Python 3.14
t-strings).Replaces Many Legacy Tools
Consolidates pip, pip-tools, pipx, poetry, pyenv, twine, and virtualenv into a single fast CLI.
Blazing Fast (Written in Rust)
10–100x faster package installation and lock resolution. Drastically shortens CI pipeline runs.
Universal Lockfile (uv.lock)
Pins the entire dependency graph across platforms and Python versions for 100% reproducible environments.
Single-File Script Execution
PEP 723 inline dependency metadata allows single-file scripts to declare and isolate their own packages.
Manages Python Interpreters
Downloads, installs, and pins Python runtimes per project automatically without system packages.
Cargo-Style Workspaces
Native workspace support for modular multi-package monorepos and microservice agent collections.
| Command | Operational Purpose |
|---|---|
uv init my-agent |
Create a new project with pyproject.toml, .python-version, and entry point |
uv add httpx pydantic |
Add runtime dependencies; automatically updates pyproject.toml and uv.lock |
uv remove httpx |
Remove a dependency cleanly from files and virtual environment |
uv sync |
Reconcile the local .venv to match uv.lock with surgical precision |
uv run pytest |
Run any command inside the virtual environment (auto-syncs environment first) |
uv lock --upgrade |
Re-resolve all dependencies and upgrade pinned package versions |
uv tree |
Inspect and visualize the full transitive dependency tree |
Golden Rule:
.venvanduv.lockare created automatically on firstuv addoruv run. You never activate the environment by hand.
# Add development tools to dev group
uv add --dev pytest ruff pytest-asyncio
# Add documentation dependencies
uv add --group docs mkdocs
# Developer workstation: sync everything
uv sync
# CI environment: fail if lockfile changes
uv sync --frozen
# Production container: exclude test/lint tools
uv sync --no-dev
# Export for legacy platforms requiring requirements.txt
uv export --format requirements-txt > requirements.txtuv.lock: Guarantee that every developer and every CI runner gets bit-for-bit identical packages.uv add and uv remove so pyproject.toml and uv.lock stay strictly synchronized.uv run executes the script inside an isolated ephemeral cache, ignoring the surrounding project’s dependencies.uv init --script agent_eval.pyuv add --script agent_eval.py anthropicuv run --with duckdb query.pyuvx and uv tool| Command | Purpose | When to Use |
|---|---|---|
uvx ruff check . |
Run CLI tool in temporary environment | One-off checks, formatting PRs, quick audits |
uv tool install ruff |
Install CLI tool globally, isolated from system Python | Everyday tools you invoke across all terminal directories |
uv tool upgrade --all |
Upgrade all globally installed CLI tools | Routine developer workstation maintenance |
uv run --with pytest pytest |
Run command with ad-hoc package injected | Testing scripts without polluting pyproject.toml |
Ruff
Extremely fast Python linter & formatter in Rust; replaces flake8, black, isort, bandit.
Type Checker
mypy or pyright for static safety checks in CI and editors.
Pytest
The industry-standard test runner with fixture and async support.
pyproject.toml[project]
name = "my-agent"
version = "0.1.0"
description = "Autonomous Agent with LangGraph & ADK"
readme = "README.md"
requires-python = ">=3.12"
dependencies = [
"pydantic>=2.7",
"httpx>=0.27.0",
]
[dependency-groups]
dev = [
"pytest>=8.0.0",
"pytest-asyncio>=0.23.0",
"ruff>=0.4.0",
"mypy>=1.10.0",
]
[tool.ruff]
line-length = 100
[tool.pytest.ini_options]
asyncio_mode = "auto"[project]: Defines package name, semver version, Python compatibility boundary, and core runtime requirements.[dependency-groups]: Standardized dev, test, and documentation groups (replaces legacy poetry extras).[tool.*]: Centralized configuration for developer tools (ruff, pytest, mypy), eliminating cluttered standalone config files.my-agent/
├── pyproject.toml # Single source of truth
├── uv.lock # Cross-platform deterministic lock
├── .python-version # Pinned Python interpreter (e.g. 3.13)
├── .env # Environment secrets (NEVER commit!)
├── .gitignore # Git exclusion rules
├── src/
│ └── my_agent/ # Primary package root
│ ├── __init__.py
│ ├── agent.py # Core agent loop / state machine
│ ├── tools.py # Exposed tools & schemas
│ ├── schemas.py # Pydantic models
│ └── config.py # Settings & env vars
└── tests/
├── conftest.py # Global fixtures & mock models
└── test_tools.py # Unit tests mirroring src
src Layout: Prevents accidental imports of un-built local directory code when running from repository root.src/my_agent/<name>.py corresponds to tests/test_<name>.py.I (isort): Ensures deterministic, clean import structures.UP (pyupgrade): Automatically updates older typing syntax to modern Python 3.12+ syntax (int | None).ASYNC: Catches catastrophic bugs in agent code, such as calling time.sleep() or requests.get() inside async functions.Team Culture: Let the formatter end style debates. Never disable a rule without an explicit
# noqacomment explaining the engineering reason.
Pick a Checker
Ratchet Up Strictness
basic mode).disallow_untyped_defs, no_implicit_optional) incrementally package by package..env or system environment.ruff check and ruff format before commits are created.Irreversible Exposure
A key committed once lives in Git commit history forever. If a secret leaks, rotate it immediately in the cloud provider console.
Goal
Create a reproducible Python agent repository from scratch with linting, typing, and tests configured.
uv init my-agent && uv python pin 3.12uv add httpx pydanticuv add --dev pytest pytest-asyncio ruff mypytool.ruff and tool.pytest.ini_options in pyproject.tomltests/test_basic.py; execute with uv run pytest.venv and rebuild completely via uv sync# pyproject.toml - Fully wired configuration
[project]
name = "my-agent"
version = "0.1.0"
description = "Autonomous Agent Core Skeleton"
readme = "README.md"
requires-python = ">=3.12"
dependencies = [
"httpx>=0.27.0",
"pydantic>=2.7.0",
]
[dependency-groups]
dev = [
"pytest>=8.0.0",
"pytest-asyncio>=0.23.0",
"ruff>=0.4.0",
"mypy>=1.10.0",
]
[tool.ruff]
line-length = 100
[tool.pytest.ini_options]
testpaths = ["tests"]
asyncio_mode = "auto"# Shared references can lead to unintended mutations
a = [1, 2]
b = a # Same list in memory, two names
b.append(3)
print(a) # [1, 2, 3]!
# Making an independent copy
import copy
c = copy.deepcopy(a)
print(a is b) # True (same identity)
print(a is c) # False (different objects)
print(a == c) # True (equal values)In Agent Frameworks
State updates in LangGraph, Strands, and ADK must return new values rather than mutating shared state in place. In-place mutations corrupt replay checkpoints, break parallel branches, and invalidate rollbacks.
def invoke_tool(
tool_name: str, # Positional or keyword
/, # Positional-only boundary
query: str, # Positional or keyword
*, # Keyword-only boundary
timeout: float = 30.0,
retries: int = 3,
**kwargs
) -> dict:
"""Execute an agent tool with strict parameter rules."""
return {"tool": tool_name, "query": query, "timeout": timeout}
# Valid call:
invoke_tool("search", "python agent", timeout=10.0)
# Invalid: invoke_tool(tool_name="search", query="...") raises TypeError*: Forces caller clarity for options (timeout, model), preventing accidental parameter swapping.# Comprehensions: eager memory allocation
squares = [n * n for n in range(10) if n % 2 == 0]
lookup = {u.id: u for u in users}
# Generator Expression: lazy, constant memory
total = sum(n * n for n in range(10**7))
# Generator Function: streams lines one at a time
def stream_dataset(path: str):
with open(path, "r", encoding="utf-8") as f:
for line in f:
yield line.rstrip()
# Python 3.12+ batching
from itertools import batched, islice
first_five = list(islice(stream_dataset("corpus.jsonl"), 5))
batches = list(batched(range(10), 3)) # [(0,1,2), (3,4,5), (6,7,8), (9,)]__iter__() and __next__().class AgentToolError(Exception):
"""Recoverable tool failure that can be fed back to LLM."""
try:
response = call_external_service()
except TimeoutError as exc:
# Preserve root traceback via 'from'
raise AgentToolError("Search provider timed out") from exc
finally:
cleanup_resources()
# Python 3.11+ ExceptionGroup for concurrent tasks
try:
...
except* AgentToolError as eg:
for e in eg.exceptions:
log.warning(f"Tool failed: {e}")raise ... from exc: Chains exceptions cleanly, giving the LLM or operator the complete causal story.except BaseException: or bare except:.KeyboardInterrupt and asyncio.CancelledError, preventing the runtime from gracefully stopping hanging agent runs.# 1. Standard protocol (__enter__ and __exit__)
with open("checkpoint.json", "r") as f:
data = json.load(f)
# 2. Generator-based with contextlib
import time
from contextlib import contextmanager
@contextmanager
def execution_timer(label: str):
start = time.perf_counter()
try:
yield
finally:
elapsed = time.perf_counter() - start
print(f"[{label}] Elapsed: {elapsed:.3f}s")
with execution_timer("Vector Search"):
results = perform_search("query")__exit__ executes even when unhandled exceptions occur.contextlib.ExitStack: Dynamically manage an arbitrary number of open resources (e.g. dynamic MCP tool servers).import functools
from collections.abc import Callable
TOOL_REGISTRY: dict[str, Callable] = {}
def agent_tool(func: Callable) -> Callable:
"""Decorator registering functions as agent tools."""
@functools.wraps(func) # Crucial: preserves name & docstring!
def wrapper(*args, **kwargs):
return func(*args, **kwargs)
TOOL_REGISTRY[func.__name__] = wrapper
return wrapper
@agent_tool
def calculate_tax(amount: float, rate: float) -> float:
"""Calculate the sales tax for an invoice."""
return amount * rate
# The registry now has calculate_tax registered with its docstring intactfunctools.wraps: Without this, calculate_tax.__name__ becomes wrapper and __doc__ is lost, breaking tool schema generation.@retry(times=3)) requires an extra nesting layer that returns the actual decorator.# Structural Pattern Matching (Python 3.10+)
match agent_event:
case {"type": "text", "content": str(text)}:
display_message(text)
case {"type": "tool_call", "name": name, "args": dict(args)}:
execute_tool(name, args)
case {"type": "error", "code": int(c)} if c >= 500:
trigger_fallback()
case _:
pass
# Modern unpacking & merging
first, *rest = message_history
config = default_config | user_overridescollections.deque(maxlen=100): Perfect for sliding conversational memory windows where old tokens drop off automatically.collections.defaultdict(list): Grouping tool calls by provider or status without boilerplate key existence checks.collections.Counter: Counting token usages and tool call frequencies.# String Formatting & Templates
name, count = "Claude", 42
print(f"{name=}, {count=}") # Debug form
from textwrap import dedent
system_prompt = dedent(f"""\
You are an AI assistant named {name}.
Active tools: {count}.
""").strip()
# Python 3.14 Preview: PEP 750 t-strings
# t"SELECT * FROM agents WHERE name = {name}"
# Allows safe template parsing without SQL injection.Enums over Magic Strings: String typos like
"tool_ues"fail silently;StopReason.tool_uesraises an immediateAttributeError.
Breaking Circular Imports
| Module | Core Purpose in AI Engineering |
|---|---|
pathlib |
Modern object-oriented path traversal: Path("data") / "prompts" / "v1.txt" |
json, tomllib |
Parse tool arguments, read configuration TOML files (Python 3.11+) |
datetime, zoneinfo |
Timezone-aware timestamps for chat logs: datetime.now(timezone.utc) |
logging |
Hierarchical, leveled logging with contextual metadata per agent run |
os, argparse |
Read environment variables (os.getenv), parse CLI arguments |
uuid, hashlib |
Generate unique run IDs (uuid4()), hash prompt cache keys (sha256) |
functools, operator |
Function memoization (@cache), reducers (operator.add in LangGraph) |
Goal
Build a lazy text processing pipeline with decorators, timing context managers, and robust error chaining.
@stage decorator that registers processing functions into a pipeline table.@contextmanager timer that logs the execution duration of each stage.PipelineError using raise ... from.defaultdict and Counter.import functools
import time
from collections import Counter
from collections.abc import Callable, Generator
from contextlib import contextmanager
class PipelineError(Exception):
"""Raised when a text processing stage fails."""
STAGES: dict[str, Callable] = {}
def stage(name: str) -> Callable:
"""Decorator registering pipeline processing steps."""
def decorator(fn: Callable) -> Callable:
@functools.wraps(fn)
def wrapper(*args, **kwargs):
return fn(*args, **kwargs)
STAGES[name] = wrapper
return wrapper
return decorator
@contextmanager
def stage_timer(label: str):
start = time.perf_counter()
try:
yield
finally:
print(f"[{label}] Elapsed: {time.perf_counter() - start:.4f}s")@stage("clean_tokens")
def stream_cleaned(
lines: Generator[str, None, None]
) -> Generator[str, None, None]:
for line in lines:
for word in line.lower().split():
cleaned = word.strip(".,!?:;\"'()[]")
if cleaned:
yield cleaned
def execute_pipeline(raw_lines: list[str]) -> Counter:
with stage_timer("token_aggregation"):
try:
line_generator = (line for line in raw_lines)
tokens = stream_cleaned(line_generator)
return Counter(tokens)
except Exception as exc:
raise PipelineError("Stage failure") from exc
# Run pipeline:
results = execute_pipeline(["Agent workflows require async.", "Async is fast."])
print(results)
# Counter({'async': 2, 'agent': 1, 'workflows': 1, ...})Frameworks Read Them
Tool schemas, state channels, and structured outputs are derived from annotations dynamically at runtime.
Checkers Read Them
Static type checkers catch interface mismatches, missing fields, and type bugs before executing expensive LLM API calls.
Humans Read Them
Annotations document strict contracts between nodes, tools, and multi-agent coordination teams.
Python 3.14 & PEP 649
Python 3.14 defers annotation evaluation until explicitly inspected. Use annotationlib.get_annotations() with Format.VALUE, Format.FORWARDREF, or Format.STRING to safely introspect types without circular import errors.
# Modern union syntax (Python 3.10+) & built-in generics
def process_messages(
messages: list[str],
user_id: int | None = None,
options: dict[str, float] | None = None
) -> tuple[int, str]:
...
# Accept abstract types, return concrete types
from collections.abc import Callable, Sequence, Mapping
def execute_pipeline(
handler: Callable[[str], dict],
inputs: Sequence[str], # Accepts list, tuple, deque
config: Mapping[str, str] # Accepts dict, custom mapping
) -> list[dict]:
return [handler(item) for item in inputs]|: Use str | None instead of legacy Optional[str].list[T], dict[K, V], set[T] instead of importing typing.List or typing.Dict.Sequence or Iterable in parameters so callers can pass lists, tuples, or generators interchangeably.Literal, Final, and Annotatedfrom typing import Literal, Final, Annotated
from operator import add
from pydantic import Field
# Literal: strictly allowed string options for routing
AgentMode = Literal["planner", "executor", "evaluator"]
# Final: constants that cannot be reassigned
MAX_SUBTASKS: Final[int] = 10
# LangGraph-style Reducer Metadata:
# Annotated[Type, ReducerFunction]
class AgentState(TypedDict):
messages: Annotated[list[str], add]
scratchpad: Annotated[list[str], add]
# Pydantic-style Field Constraint:
Confidence = Annotated[float, Field(ge=0.0, le=1.0)]Literal: Drives router nodes in LangGraph and model decision points.Final: Pins safety constants (turn limits, token thresholds).Annotated[T, metadata]: Attaches runtime instructions directly to the type without changing static checking. LangGraph inspects the second argument to know how to merge partial updates (operator.add).TypedDict: Lightweight State Modelingfrom typing import TypedDict, NotRequired, Required
class ChatMessage(TypedDict):
role: str
content: str
tool_call_id: NotRequired[str] # Optional key
# Python 3.13 ReadOnly support
from typing import ReadOnly
class SessionState(TypedDict, total=False):
session_id: Required[str] # Must be present
messages: list[ChatMessage]
api_key: ReadOnly[str] # Immutable key
msg: ChatMessage = {"role": "user", "content": "hello"}TypedDictdict. No instantiation latency.NamedTuplefrom dataclasses import dataclass, field, replace
@dataclass(frozen=True, slots=True, kw_only=True)
class RuntimeContext:
user_id: str
agent_name: str = "assistant"
tags: list[str] = field(default_factory=list)
def __post_init__(self):
if not self.user_id:
raise ValueError("user_id cannot be empty")
ctx1 = RuntimeContext(user_id="usr_123")
# Functional non-mutating update:
ctx2 = replace(ctx1, agent_name="specialist")slots=True: Drastically reduces memory consumption and speeds up attribute lookups.frozen=True: Enforces immutability; instances are hashable and thread-safe.default_factory: Always use field(default_factory=list) to avoid the shared mutable container trap.replace(): Creates a modified clone without mutating the original.Protocol vs ABC: Describing BehaviorProtocol)from typing import Protocol
class ModelClient(Protocol):
async def generate(self, prompt: str) -> str:
...
# Conforms automatically WITHOUT inheritance!
class FakeModelClient:
async def generate(self, prompt: str) -> str:
return "mock response"
def run_agent(client: ModelClient):
...
run_agent(FakeModelClient()) # Type-checks cleanly!Rule of Thumb: Use
Protocolfor test seams, fakes, and client swapping. UseABCwhen sharing concrete helper code among derived plugin classes.
# Modern PEP 695 Syntax (Python 3.12+)
def first_item[T](items: list[T]) -> T:
return items[0]
class StateStore[TState, TConfig]:
def __init__(self, state: TState, config: TConfig) -> None:
self.state = state
self.config = config
# Type Aliases
type AgentResult[T] = tuple[bool, T, dict[str, int]]
# Bounded Generics
class Runner[T: ModelClient]:
def __init__(self, client: T) -> None:
self.client = clientTypeVar: Replaces verbose T = TypeVar("T") declarations with clean syntax parameters [T].type Keyword: First-class statement for defining parameterized type aliases.Runtime[Context], Case[str, str], and Channel[T].ParamSpec & Typed Decoratorsfrom collections.abc import Callable
from functools import wraps
def with_retry[**P, R](attempts: int = 3):
def decorator(fn: Callable[P, R]) -> Callable[P, R]:
@wraps(fn)
def wrapper(*args: P.args, **kwargs: P.kwargs) -> R:
for i in range(attempts):
try:
return fn(*args, **kwargs)
except Exception:
if i == attempts - 1:
raise
raise RuntimeError("Exhausted")
return wrapper
return decoratorParamSpec is CrucialParamSpec, a decorator would type the function as Callable[..., Any], destroying autocomplete and runtime tool schema inspection!from typing import Literal, TypedDict, assert_never
class TextChunk(TypedDict):
kind: Literal["text"]
text: str
class ToolChunk(TypedDict):
kind: Literal["tool_call"]
name: str
args: dict
StreamChunk = TextChunk | ToolChunk
def handle_chunk(chunk: StreamChunk) -> None:
match chunk["kind"]:
case "text":
print(chunk["text"]) # Narrowed to TextChunk
case "tool_call":
call_tool(chunk["name"]) # Narrowed to ToolChunk
case _ as unreachable:
assert_never(unreachable) # Compile-time exhaustiveness!Literal field (kind) allows the type checker to narrow branches cleanly.assert_never(): If a new event type (e.g. ErrorChunk) is added to StreamChunk but forgotten in the match statement, the type checker immediately raises a static build error.import inspect
from typing import get_type_hints
def generate_tool_schema(fn) -> dict:
"""Inspect a Python function and construct JSON Schema."""
hints = get_type_hints(fn)
sig = inspect.signature(fn)
parameters = {}
for param_name, param in sig.parameters.items():
param_type = hints.get(param_name, str).__name__
parameters[param_name] = {
"type": param_type,
"default": None if param.default is inspect._empty else param.default
}
return {
"name": fn.__name__,
"description": inspect.getdoc(fn) or "",
"parameters": parameters
}@tool.Common Traps
Any Leaks: One Any silently infects everything downstream. Use object if uncertain.x: int | None still requires an argument unless declared as = None.Incremental Strictness Plan
disallow_untyped_defs for core modules.--strict across the whole repository in CI.Goal
Model an agent workflow with precise types, then generate JSON tool schemas directly from Python function signatures.
AgentState as a TypedDict with NotRequired keys.Annotated reducer and apply it in a custom state merge function.TypedDict structures.match ... case and assert_never.ToolRegistry[T] using Python 3.12+ PEP 695 syntax.import inspect
from operator import add
from typing import Annotated, Literal, NotRequired, TypedDict, assert_never, get_type_hints
# 1. Typed State with Annotated reducer
class AgentState(TypedDict):
turn: int
messages: Annotated[list[str], add]
scratchpad: NotRequired[dict[str, str]]
# 2. Discriminated union of events
class TextEvent(TypedDict):
kind: Literal["text"]
text: str
class ToolEvent(TypedDict):
kind: Literal["tool"]
name: str
type AgentEvent = TextEvent | ToolEvent
def dispatch_event(event: AgentEvent) -> str:
match event["kind"]:
case "text":
return f"Emitted: {event['text']}"
case "tool":
return f"Called: {event['name']}"
case _ as unreachable:
assert_never(unreachable)# 3. PEP 695 Generic Registry
class ToolRegistry[T]:
def __init__(self) -> None:
self._registry: dict[str, T] = {}
def register(self, name: str, item: T) -> None:
self._registry[name] = item
def get(self, name: str) -> T:
return self._registry[name]
# 4. JSON Schema Extractor from Signature
def extract_schema(fn) -> dict:
hints = get_type_hints(fn)
sig = inspect.signature(fn)
props = {
name: {"type": hints.get(name, str).__name__}
for name in sig.parameters
}
return {
"name": fn.__name__,
"description": inspect.getdoc(fn) or "",
"parameters": props,
}Untrusted Input
Model outputs, external webhooks, and tool responses are untrusted raw data until strictly validated.
Schemas for Free
Pydantic automatically emits OpenAPI / JSON Schema specifications required by tool-calling models.
Everywhere in Agents
Strands structured outputs, ADK input/output contracts, FastAPI endpoints, and application settings.
Rule of Thumb: Types document intent, Pydantic enforces reality. Validate once at the system boundary, then pass typed, guaranteed objects inward.
BaseModel Essentialsfrom pydantic import BaseModel, ValidationError
class FlightBooking(BaseModel):
flight_code: str
passengers: int = 1
tags: list[str] = []
# Automatic type coercion (lax mode by default)
b = FlightBooking(flight_code="BA249", passengers="3")
print(b.passengers) # int: 3
# Serialization
dumped_dict = b.model_dump()
json_string = b.model_dump_json()
# Deserialization & validation
try:
FlightBooking.model_validate({"passengers": 0})
except ValidationError as e:
print(e.errors()) # Location & message for missing fields"3" \(\rightarrow\) 3). Use strict=True to reject coercion.ValidationError: Pinpoints the exact path (loc), invalid input value, and error type.model_json_schema(): Extracts the exact JSON Schema for model tool definitions.from pydantic import BaseModel, Field, field_validator, model_validator
class ToolOrder(BaseModel):
quantity: int = Field(ge=1, le=100, description="Items to order")
sku: str = Field(pattern=r"^[A-Z]{3}-\d{4}$")
min_price: float
max_price: float
@field_validator("sku")
@classmethod
def normalize_sku(cls, v: str) -> str:
return v.strip().upper()
@model_validator(mode="after")
def validate_price_range(self):
if self.min_price > self.max_price:
raise ValueError("min_price cannot exceed max_price")
return selfField(description=...): Injected directly into the model’s prompt, explaining semantics and value boundaries to the LLM.@field_validator: Cleans or checks single fields.@model_validator(mode="after"): Validates cross-field constraints after initial field parsing.from typing import Literal, Annotated
from pydantic import BaseModel, Field
class TextContent(BaseModel):
kind: Literal["text"]
text: str
class ImageContent(BaseModel):
kind: Literal["image"]
image_url: str
# Tagged / Discriminated Union
ContentPart = Annotated[TextContent | ImageContent, Field(discriminator="kind")]
class AgentMessage(BaseModel):
role: str
parts: list[ContentPart]
# Validates and picks the exact subclass instantly:
msg = AgentMessage.model_validate({
"role": "user",
"parts": [{"kind": "text", "text": "Analyze this chart"}]
})kind) instead of trial-and-error parsing across every union member.image_url is missing in an image part, Pydantic reports image_url missing rather than confusing union union-exhaustion errors.model_rebuild() for models referencing themselves (e.g. tree structures).from pydantic import SecretStr
from pydantic_settings import BaseSettings, SettingsConfigDict
class AgentSettings(BaseSettings):
model_config = SettingsConfigDict(
env_file=".env",
env_file_encoding="utf-8",
extra="forbid"
)
anthropic_api_key: SecretStr
model_name: str = "claude-3-5-sonnet"
max_tokens: int = 4096
timeout_seconds: float = 30.0
settings = AgentSettings() # Automatically reads from .env and os.environSecretStr: Masks API keys in logs and tracebacks (**********), preventing credential leaks.TypeAdapter & LLM Self-Correctionfrom pydantic import BaseModel, TypeAdapter, ValidationError
class ExtractedEntity(BaseModel):
name: str
category: str
confidence: float
# Validate top-level list without wrapping in a container model
adapter = TypeAdapter(list[ExtractedEntity])
# LLM Error-Feedback Retry Loop
for attempt in range(3):
raw_json = call_llm(extraction_prompt)
try:
entities = adapter.validate_json(raw_json)
break
except ValidationError as exc:
extraction_prompt += f"\nYour JSON was invalid:\n{exc.errors()}\nPlease fix and re-emit."
else:
raise RuntimeError("LLM failed to output valid schema after 3 retries")TypeAdapter: Validates native types like list[Model], dict[str, int], or TypedDict directly.validate_json(): Parses and validates in a single Rust-accelerated pass.TypedDict, Dataclass, or Pydantic?| Characteristic | TypedDict |
Dataclass | Pydantic Model |
|---|---|---|---|
| Runtime Validation | No | No | Yes |
| Instance Type | Plain dict |
Class instance | Class instance |
| Execution Speed | Fastest (\(O(1)\)) | Fast | Slower (validation cost) |
| Schema Output | Via adapter | Via adapter | Built-in native |
| Best Used For | High-throughput graph state | Internal config & context | System boundaries & LLM I/O |
Architecture Decision: LangGraph internal graph state defaults to
TypedDictfor raw speed; Strands and ADK lean on Pydantic for tool validation and structured output. Choose based on where the data lives!
Goal
Define, validate, and document the complete data exchange contract for an autonomous agent.
ToolRequest and ToolResponse with numeric and string constraints using Field.@model_validator enforcing cross-field logical consistency.model_json_schema().Settings class using pydantic-settings to load API keys securely from .env.TypeAdapter.from typing import Annotated, Literal
from pydantic import BaseModel, Field, SecretStr, TypeAdapter, ValidationError, model_validator
from pydantic_settings import BaseSettings, SettingsConfigDict
# 1. Discriminated union of multimodal parts
class TextPart(BaseModel):
kind: Literal["text"]
content: str
class ToolPart(BaseModel):
kind: Literal["tool_call"]
tool_name: str
args: dict
Part = Annotated[TextPart | ToolPart, Field(discriminator="kind")]
# 2. Tool Contract with validation
class ToolRequest(BaseModel):
request_id: str
parts: list[Part]
max_retries: int = Field(ge=1, le=5, default=3)
@model_validator(mode="after")
def ensure_non_empty(self):
if not self.parts:
raise ValueError("parts cannot be empty")
return self# 3. Settings configuration with SecretStr
class AgentConfig(BaseSettings):
model_config = SettingsConfigDict(env_file=".env", extra="ignore")
api_key: SecretStr = SecretStr("mock_secret_key")
environment: str = "production"
# 4. TypeAdapter self-healing validation loop
adapter = TypeAdapter(ToolRequest)
def robust_parse(json_str: str) -> ToolRequest:
try:
return adapter.validate_json(json_str)
except ValidationError as err:
# Structured errors to send back to model
error_hints = [f"{e['loc']}: {e['msg']}" for e in err.errors()]
raise RuntimeError(f"Model error hints: {error_hints}") from err
# Example valid JSON:
payload = '{"request_id": "r1", "parts": [{"kind": "text", "content": "run audit"}]}'
validated = robust_parse(payload)
print(validated.model_dump())| Model | Ideal Workload | Context Switching | Shared State | Task Cost |
|---|---|---|---|---|
asyncio |
Many waiting I/O calls | Cooperative (at await) |
Single thread, shared heap | Ultra-low (microseconds) |
| Threads | Blocking legacy libraries | Preemptive (OS-driven) | Shared memory (needs locks) | Moderate (\(~8\text{MB}\) stack) |
| Processes | CPU-heavy compute / math | Preemptive (OS-driven) | Isolated memory (pipes/IPC) | High (process fork) |
Agents Wait Most of the Time
Model inference, vector search, web browsing, and tool APIs spend 95%+ of latency waiting for network sockets. asyncio is the ideal paradigm.
The Single Event Loop
Tasks take turns cooperatively. If any single function executes synchronous blocking code, every other agent task on that loop freezes completely.
awaitimport asyncio
async def fetch_embeddings(doc_id: int) -> list[float]:
await asyncio.sleep(0.1) # Yields execution back to loop
return [0.1 * doc_id, 0.2, 0.3]
async def main():
# 1. Coroutine object: created, not started
coro = fetch_embeddings(1)
# 2. Await: runs to completion inline
result1 = await coro
# 3. create_task: scheduled concurrently in background
task = asyncio.create_task(fetch_embeddings(2))
# Loop can do other work here...
result2 = await task
print(result1, result2)
asyncio.run(main()) # Bootstraps the event loopasync def: Calling the function returns a coroutine object immediately without executing its body.await: Suspends the current coroutine until the awaited object produces a value, freeing the loop to run other work.create_task(): Wraps a coroutine in an asyncio.Task and adds it to the loop’s ready queue immediately.gather vs TaskGroupasyncio.gather (Legacy Fan-Out)TaskGroup (Python 3.11+ Structured Concurrency)t1 crashes, t2 is automatically cancelled.ExceptionGroup.import asyncio
async def background_worker():
try:
while True:
await asyncio.sleep(1)
except asyncio.CancelledError:
# Mandatory: perform async cleanup
await flush_audit_logs()
raise # NEVER swallow CancelledError!
finally:
# Runs on both clean exit and cancellation
close_local_resources()
async def execute_with_timeout():
try:
async with asyncio.timeout(5.0): # Python 3.11+
await background_worker()
except TimeoutError:
print("Agent tool exceeded 5-second SLA")asyncio.timeout(): Context manager that cancels the inner block cleanly upon expiration.CancelledError: Always re-raise CancelledError. Swallowing it prevents the runtime from terminating cancelled agent requests.finally blocks.import asyncio
# 1. Semaphore: enforce LLM rate limits
rate_limiter = asyncio.Semaphore(5)
async def rate_limited_call(url: str):
async with rate_limiter: # At most 5 concurrent requests
return await http_client.get(url)
# 2. Queue: decoupled producer-consumer with backpressure
work_queue: asyncio.Queue[str] = asyncio.Queue(maxsize=50)
async def producer():
for item in dataset:
await work_queue.put(item) # Suspends if queue is full!
async def worker():
while True:
item = await work_queue.get()
await process(item)
work_queue.task_done()asyncio.Semaphore: Prevents hitting 429 Too Many Requests errors when calling external LLMs in parallel.asyncio.Queue: Provides natural backpressure so producers don’t flood system memory.asyncio.Lock & Event: Thread-safe mutual exclusion and flag-based synchronization between agent nodes.import asyncio
from contextlib import aclosing
async def stream_tokens(prompt: str):
"""Async generator yielding LLM tokens as they arrive."""
for word in prompt.split():
await asyncio.sleep(0.05) # Simulate network chunk
yield word + " "
async def render_chat():
# Consume tokens in real time
async for token in stream_tokens("Autonomous agents reason in steps"):
print(token, end="", flush=True)
# Safe early termination with aclosing
async def partial_consume():
async with aclosing(stream_tokens("hello world")) as stream:
async for token in stream:
if "hello" in token:
break # Stream generator properly closed!yield inside an async def function creates an AsyncIterator.async for: Consumes chunks one at a time as the remote model yields them.import httpx
from contextlib import asynccontextmanager
@asynccontextmanager
async def agent_session():
# Initialize connection pool once
limits = httpx.Limits(max_keepalive_connections=20, max_connections=50)
async with httpx.AsyncClient(limits=limits, timeout=15.0) as client:
yield client
# Client closed and connections drained cleanly here
async def query_cluster(urls: list[str]):
async with agent_session() as client:
tasks = [client.get(u) for u in urls]
return await asyncio.gather(*tasks)httpx.AsyncClient per request incurs TLS handshake and DNS overhead on every call.import asyncio
import time
from concurrent.futures import ProcessPoolExecutor
# FATAL ANTI-PATTERN (Freezes all concurrent agents!):
# time.sleep(2)
# requests.get("https://api.openai.com")
# CORRECT PATTERN: Offload blocking I/O to thread pool
def blocking_disk_read(path: str) -> str:
with open(path) as f: return f.read()
content = await asyncio.to_thread(blocking_disk_read, "large.txt")
# Offload CPU-heavy tokenization / math to ProcessPool
process_pool = ProcessPoolExecutor()
loop = asyncio.get_running_loop()
embeddings = await loop.run_in_executor(
process_pool, cpu_intensive_math, matrix
)Symptoms: WebSocket heartbeats fail, token streams stutter, health checks time out.
Detection:
Logs warnings whenever a callback blocks the loop for \(>100\text{ms}\).
Static Prevention: Ruff ASYNC rules flag requests or time.sleep in async code.
| Scenario | Recommended Architectural Approach |
|---|---|
| Sync CLI script calls async agent | asyncio.run(main()) once at the application entry point |
| Already inside running loop (FastAPI, Jupyter) | await coro directly; calling asyncio.run() raises RuntimeError |
| Async agent must invoke blocking sync code | await asyncio.to_thread(sync_func, *args) |
| Framework offers sync and async methods | Use ainvoke() / astream() in async web apps; use invoke() in batch scripts |
| Synchronous SDK wrapper around async core | Place a single asyncio.run() at the outermost boundary |
Warning: Never nest
asyncio.run()calls. Avoidnest_asyncioin production architectures as it masks broken event loop designs.
contextvars: State That Follows a Taskimport contextvars
import logging
# Define context variable with default
run_id_var: contextvars.ContextVar[str] = contextvars.ContextVar(
"run_id", default="unknown"
)
async def tool_step():
# Automatically inherits the task's context without passing arguments!
current_run = run_id_var.get()
logging.info(f"[{current_run}] Executing tool")
async def handle_agent_turn(turn_id: str):
run_id_var.set(turn_id)
await tool_step()
# Tasks receive an isolated copy of contextvars:
# Updating run_id in Task A does NOT overwrite Task B!contextvarsrun_id, trace_id, and tenant_id across deeply nested agent tool functions.create_task, but modifications stay scoped to that branch.The GIL
In standard CPython, the Global Interpreter Lock ensures only one thread executes Python bytecode at a time. Threads still accelerate blocking C-level I/O.
Python 3.14 Free-Threaded
Official support for builds without the GIL (PEP 703). Multi-threaded CPU algorithms achieve true hardware parallel speedups.
Subinterpreters (PEP 734)
Run multiple independent Python runtimes inside a single process with isolated GILs and memory spaces.
import pytest
from unittest.mock import AsyncMock
@pytest.mark.asyncio
async def test_agent_fetch():
# Mocking async dependencies
mock_client = AsyncMock()
mock_client.get.return_value.json.return_value = {"answer": 42}
# Time-box the test to guarantee CI never hangs
async with asyncio.timeout(2.0):
result = await run_agent_query(mock_client, "test")
assert result == 42
mock_client.get.assert_awaited_once()
# Testing async generator streams
async def test_streaming():
chunks = [c async for c in stream_tokens("hi there")]
assert chunks == ["hi ", "there "]pytest-asyncio: Enables async def test_* functions.pytest.ini: Configure asyncio_mode = "auto" to avoid manually decorating every test.AsyncMock: Accurately mimics awaitable functions and tracks assert_awaited_once().asyncio.timeout to prevent hung sockets from freezing test runs.| Dangerous Anti-Pattern | Why It Breaks | Clean Idiomatic Solution |
|---|---|---|
Forgetting await |
Coroutine object created but never executed | Add await or wrap in asyncio.create_task |
time.sleep() in async code |
Freezes all running tasks on the entire thread | Use await asyncio.sleep() or asyncio.to_thread |
Swallowing CancelledError |
Task refuses to die during shutdowns/timeouts | Clean up in finally and always re-raise |
Fire-and-forget create_task |
Garbage collector may destroy task mid-flight | Store reference in a set or use TaskGroup |
Unbounded asyncio.gather |
Exhausts file descriptors or triggers API 429s | Throttle with asyncio.Semaphore |
asyncio.run in existing loop |
Crashes with RuntimeError: loop is running |
Await the coroutine directly |
| New HTTP client per request | No socket reuse; connection exhaust | Maintain a single shared httpx.AsyncClient |
Goal
Build a high-performance concurrent fetcher with rate limiting, timeouts, cancellation, and streaming output.
asyncio.gather bounded by a Semaphore(5).asyncio.timeout(1.0).asyncio.TaskGroup and handle ExceptionGroup.yield chunks as they complete).asyncio.to_thread.import asyncio
from contextvars import ContextVar
from collections.abc import AsyncGenerator
trace_id_var: ContextVar[str] = ContextVar("trace_id", default="anon")
async def fetch_endpoint(url: str, sem: asyncio.Semaphore) -> str:
async with sem:
try:
async with asyncio.timeout(1.5):
await asyncio.sleep(0.1) # Simulate network latency
return f"[{trace_id_var.get()}] {url} -> 200 OK"
except asyncio.CancelledError:
print(f"Aborting {url} safely...")
raise
async def stream_cancellable_fetches(
urls: list[str]
) -> AsyncGenerator[str, None]:
sem = asyncio.Semaphore(5)
async with asyncio.TaskGroup() as tg:
tasks = [tg.create_task(fetch_endpoint(u, sem)) for u in urls]
for t in tasks:
yield t.result()async def main():
trace_id_var.set("trace_req_8812")
urls = [f"https://agent.api/items/{i}" for i in range(12)]
try:
async with asyncio.timeout(3.0):
async for item in stream_cancellable_fetches(urls):
print(f"Streamed: {item}")
except TimeoutError:
print("Batch fetch timed out cleanly after 3.0s")
# Offload CPU-heavy file save to worker thread
def save_checkpoint(data: str):
with open("last_fetch.txt", "w") as f:
f.write(data)
await asyncio.to_thread(save_checkpoint, "Run completed successfully")
# asyncio.run(main())from collections import defaultdict
from collections.abc import Callable
class AgentEventBus:
def __init__(self):
self._listeners: dict[str, list[Callable]] = defaultdict(list)
def on(self, event_name: str):
"""Decorator to register async event listeners."""
def decorator(fn: Callable):
self._listeners[event_name].append(fn)
return fn
return decorator
async def emit(self, event_name: str, **payload):
for fn in self._listeners[event_name]:
await fn(**payload)
bus = AgentEventBus()
@bus.on("before_tool_call")
async def audit_tool(tool_name: str, args: dict):
print(f"[AUDIT] Invoking {tool_name} with {args}")on_llm_start, on_tool_end).class BaseToolPlugin:
registry: dict[str, type["BaseToolPlugin"]] = {}
name: str
def __init_subclass__(cls, **kwargs):
super().__init_subclass__(**kwargs)
if hasattr(cls, "name"):
cls.registry[cls.name] = cls
async def run(self, query: str) -> str:
raise NotImplementedError
# Auto-registers into BaseToolPlugin.registry upon class definition!
class WeatherPlugin(BaseToolPlugin):
name = "weather"
async def run(self, query: str) -> str:
return f"Weather for {query}: Sunny, 72F"
tool_instance = BaseToolPlugin.registry["weather"]()__init_subclass__: Executes automatically whenever a child class is declared, eliminating manual registration calls.pyproject.toml [project.entry-points].from dataclasses import dataclass
from typing import Protocol
class LLMProvider(Protocol):
async def complete(self, prompt: str) -> str: ...
class MemoryStore(Protocol):
def get_summary(self, user_id: str) -> str: ...
@dataclass(frozen=True)
class ExecutionDeps:
llm: LLMProvider
memory: MemoryStore
async def execute_agent_step(user_query: str, deps: ExecutionDeps) -> str:
summary = deps.memory.get_summary("u123")
prompt = f"Context: {summary}\nQuery: {user_query}"
return await deps.llm.complete(prompt)
# In unit tests: pass FakeLLM() and FakeMemory() cleanly!openai.OpenAI).RuntimeContext, ToolContext, and invocation state containers.from operator import add
# Reducer Table: mapping channel names to merge strategies
REDUCERS = {
"messages": add, # Concatenates lists: list + list
"step_count": lambda a, b: a + b
}
def apply_state_update(current_state: dict, partial_update: dict) -> dict:
new_state = dict(current_state) # Shallow copy state container
for key, value in partial_update.items():
if key in REDUCERS:
new_state[key] = REDUCERS[key](current_state.get(key, []), value)
else:
new_state[key] = value # Overwrite channel
return new_state
s0 = {"messages": ["hello"], "step_count": 0}
s1 = apply_state_update(s0, {"messages": ["agent response"], "step_count": 1})
assert s0["messages"] == ["hello"] # Original state untouched!from graphlib import TopologicalSorter
# Graph dependency definition: child -> set of prerequisites
graph_dependencies = {
"synthesize_report": {"search_web", "query_database"},
"query_database": {"parse_intent"},
"search_web": {"parse_intent"},
"parse_intent": set()
}
# Determine deterministic execution order (std library!)
sorter = TopologicalSorter(graph_dependencies)
execution_order = list(sorter.static_order())
print(execution_order)
# Output: ['parse_intent', 'query_database', 'search_web', 'synthesize_report']
# Dict-based route dispatching instead of huge if/elif blocks:
ROUTER = {"search": handle_search, "db": handle_db, "end": handle_end}graphlib.TopologicalSorter: Built into Python 3.9+; resolves execution DAGs and detects circular cycles automatically.import asyncio
import random
async def retry_with_exponential_backoff(
fn, *, attempts: int = 4, base_delay: float = 0.5
):
for attempt in range(attempts):
try:
return await fn()
except (TimeoutError, ConnectionError) as exc:
if attempt == attempts - 1:
raise
# Exponential backoff + Full Jitter
delay = base_delay * (2 ** attempt) * random.uniform(0.5, 1.5)
await asyncio.sleep(delay)TypeError, ValidationError).run_id tokens with API calls so repeated attempts don’t double-charge credit cards or create duplicate records.Goal
Write a lightweight asynchronous agent runtime featuring typed state, functional reducers, lifecycle hooks, and retries.
AgentState TypedDict with reducer metadata.graphlib.TopologicalSorter.before_node and after_node events via an async event bus.import asyncio
from graphlib import TopologicalSorter
from operator import add
from collections.abc import Callable
class MiniGraphEngine:
def __init__(self):
self._nodes: dict[str, Callable] = {}
self._deps: dict[str, set[str]] = {}
self._reducers = {"messages": add}
self._hooks: dict[str, list[Callable]] = {"before": [], "after": []}
def add_node(self, name: str, fn: Callable, depends_on: set[str] | None = None):
self._nodes[name] = fn
self._deps[name] = depends_on or set()
def register_hook(self, stage: str, fn: Callable):
self._hooks[stage].append(fn)
def _apply_reducers(self, state: dict, update: dict) -> dict:
new_state = dict(state)
for k, v in update.items():
if k in self._reducers:
new_state[k] = self._reducers[k](state.get(k, []), v)
else:
new_state[k] = v
return new_state
async def execute(self, initial_state: dict) -> dict:
sorter = TopologicalSorter(self._deps)
state = dict(initial_state)
for node_name in sorter.static_order():
for hook in self._hooks["before"]:
await hook(node_name, state)
update = await self._nodes[node_name](state)
state = self._apply_reducers(state, update)
for hook in self._hooks["after"]:
await hook(node_name, state)
return state# Execution Pipeline Demo
engine = MiniGraphEngine()
async def log_step(name: str, state: dict):
print(f"[HOOK] Executing node: {name}")
engine.register_hook("before", log_step)
# Define async graph nodes
async def fetch_step(s):
return {"messages": ["fetched doc"]}
async def summarize_step(s):
return {"summary": "Document summarized"}
engine.add_node("fetch", fetch_step)
engine.add_node("summarize", summarize_step, depends_on={"fetch"})
async def main():
final = await engine.execute({"messages": []})
print(final)
# asyncio.run(main())import httpx
# Synchronous one-off call
resp = httpx.get("https://api.openai.com/v1/models", timeout=10.0)
resp.raise_for_status()
models = resp.json()
# Asynchronous streaming (SSE / Server-Sent Events)
async def stream_agent_events(url: str):
async with httpx.AsyncClient(timeout=30.0) as client:
async with client.stream("GET", url) as stream_resp:
stream_resp.raise_for_status()
async for line in stream_resp.aiter_lines():
if line.startswith("data: "):
yield line[6:]httpx default timeout can hang indefinitely on network partitions. Always pass explicit float or httpx.Timeout configurations.raise_for_status(): Immediately converts HTTP 4xx/5xx responses into catchable Python exceptions.aiter_lines(): Reads SSE streams line-by-line without buffering full gigabyte responses into RAM.import json
import datetime as dt
from dataclasses import asdict
# Only native JSON types serialize out-of-the-box:
# dict, list, str, int, float, bool, None
json.dumps({"name": "Agent", "step": 1}) # OK
# Non-primitive trap:
# json.dumps({"timestamp": dt.datetime.now()}) # TypeError!
# Dataclass serialization
json.dumps(asdict(runtime_context))
# Pydantic handles dates, UUIDs, and enums automatically!
json_str = state_model.model_dump_json()
# Always verify round-trip serialization:
assert json.loads(json_str) is not None.value.validate(loads(dumps(state))) == state in unit tests.import os
import tempfile
import tomllib
from pathlib import Path
# Modern object-oriented paths
base_dir = Path(__file__).resolve().parent
config_path = base_dir / "config" / "agent.toml"
# Atomic file writes (prevents corrupted half-written files)
def write_atomic(target_path: Path, content: str):
target_path.parent.mkdir(parents=True, exist_ok=True)
with tempfile.NamedTemporaryFile("w", dir=target_path.parent, delete=False) as f:
f.write(content)
temp_name = f.name
os.replace(temp_name, target_path) # Atomic POSIX swap!pathlib.Path: Use / operator instead of fragile os.path.join().../../etc/passwd).import logging
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [%(levelname)s] %(name)s (run_id=%(run_id)s): %(message)s"
)
logger = logging.getLogger("agent.tools.web")
# Injecting correlation metadata
logger.info(
"Fetching document",
extra={"run_id": "run_98234"}
)
# Tune individual external libraries without noise:
logging.getLogger("httpx").setLevel(logging.WARNING)
logging.getLogger("httpcore").setLevel(logging.WARNING)my_agent.tools) lets you silence noisy third-party libraries while enabling verbose debug logs for your code.run_id to log records to trace an agent turn across distributed workers.import sqlite3
# SQLite with automatic transaction context
with sqlite3.connect("agent_memory.db") as conn:
cursor = conn.cursor()
cursor.execute("""
CREATE TABLE IF NOT EXISTS checkpoints (
run_id TEXT PRIMARY KEY,
state TEXT NOT NULL,
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
)
""")
# Parameterized query (PREVENTS SQL INJECTION)
cursor.execute(
"INSERT OR REPLACE INTO checkpoints (run_id, state) VALUES (?, ?)",
("run_123", serialized_state)
)
# Automatically committed upon exit of with block!with statement automatically commits on success and performs a rollback if an unhandled exception occurs.f"SELECT * WHERE id={id}"). Always use query parameters ? or %s.psycopg_pool or asyncpg).Goal
Fetch data concurrently across multiple endpoints, validate responses with Pydantic, persist atomically, and trace every step.
httpx.AsyncClient wrapper with timeouts and exponential retry backoff.asyncio.Semaphore(5).run_id ContextVar onto every emitted log message.import asyncio
import logging
import sqlite3
from contextvars import ContextVar
from pydantic import BaseModel, Field
run_id_var: ContextVar[str] = ContextVar("run_id", default="init")
class DocumentRecord(BaseModel):
doc_id: str
content: str
score: float = Field(ge=0.0, le=1.0)
def init_db(db_path: str = "agent_records.db"):
with sqlite3.connect(db_path) as conn:
conn.execute("""
CREATE TABLE IF NOT EXISTS docs (
doc_id TEXT PRIMARY KEY,
content TEXT,
score REAL,
run_id TEXT
)
""")
def persist_doc(doc: DocumentRecord, db_path: str = "agent_records.db"):
with sqlite3.connect(db_path) as conn:
conn.execute(
"INSERT OR REPLACE INTO docs VALUES (?, ?, ?, ?)",
(doc.doc_id, doc.content, doc.score, run_id_var.get())
)async def fetch_and_save(doc_id: str, sem: asyncio.Semaphore):
async with sem:
# Simulate network fetch & Pydantic validation
await asyncio.sleep(0.05)
doc = DocumentRecord(
doc_id=doc_id,
content=f"Content for {doc_id}",
score=0.92
)
persist_doc(doc)
logging.info(
f"Stored {doc_id}",
extra={"run_id": run_id_var.get()}
)
async def main():
init_db()
run_id_var.set("batch_2026_09")
sem = asyncio.Semaphore(5)
async with asyncio.TaskGroup() as tg:
for i in range(10):
tg.create_task(fetch_and_save(f"doc_{i}", sem))
# asyncio.run(main())pytest Essentialsimport pytest
from pathlib import Path
# Fixtures for reusable setup and teardown
@pytest.fixture
def temp_workspace(tmp_path: Path):
data_dir = tmp_path / "agent_data"
data_dir.mkdir()
yield data_dir
# Cleanup happens automatically after test!
# Parametrized tests: one test logic, many edge cases
@pytest.mark.parametrize("query,expected_category", [
("Book flight to Tokyo", "travel"),
("Reset my password", "it_support"),
("Account balance inquiry", "billing"),
])
def test_intent_classification(query, expected_category):
assert classify_intent(query) == expected_category
def test_validation_error_raised():
with pytest.raises(ValueError, match="API key required"):
initialize_agent(api_key="")assert Statements: Pytest introspects the bytecode, printing detailed left/right diffs on failure without assertEquals boilerplate.tmp_path, database connections, mock clients).@pytest.mark.parametrize: Easily test exhaustive permutations and edge cases.from unittest.mock import AsyncMock
# Fakes over Mocks: small working fake class
class FakeLLMClient:
def __init__(self, canned_response: str = "ok"):
self.canned = canned_response
self.call_history: list[str] = []
async def complete(self, prompt: str) -> str:
self.call_history.append(prompt)
return self.canned
# Testing async agent workflow with fake:
async def test_agent_with_fake():
fake = FakeLLMClient("tool_call: web_search")
result = await run_agent_loop(client=fake)
assert len(fake.call_history) == 1
# Monkeypatching environment variables safely
def test_environment_loading(monkeypatch):
monkeypatch.setenv("ANTHROPIC_API_KEY", "sk-test-key-12345")
settings = load_settings()
assert settings.api_key.get_secret_value() == "sk-test-key-12345"AsyncMock: Accurately replicates awaitables and tracks .assert_awaited_once().monkeypatch: Automatically reverts modified environment variables or module attributes after each test.Deterministic Fakes
Return fixed responses to verify state routing, JSON extraction, and error retries completely offline.
Golden Snapshot Files
Snapshot generated prompt strings and expected tool schemas to catch unintended regressions in git diffs.
Property-Based Tests
Use Hypothesis to generate random inputs, verifying mathematical properties like reducer associativity.
Mark Slow Tests
Mark tests calling real remote models with @pytest.mark.slow and run them only on scheduled pipelines.
Evals vs Unit Tests
# .github/workflows/ci.yml
name: CI Pipeline
on: [push, pull_request]
jobs:
quality-gate:
runs-on: ubuntu-latest
strategy:
matrix:
python-version: ["3.12", "3.13"]
steps:
- uses: actions/checkout@v4
- uses: astral-sh/setup-uv@v5
- name: Install dependencies
run: uv sync --frozen --python ${{ matrix.python-version }}
- name: Code formatting & linting
run: uv run ruff check && uv run ruff format --check
- name: Static type checking
run: uv run mypy src
- name: Execute test suite
run: uv run pytest --cov=src --cov-report=term-missinguv sync --frozen: Guarantees CI builds never modify uv.lock.ruff format)ruff check)mypy)Goal
Bring the Module 6 agent graph runtime under a comprehensive test suite and continuous integration pipeline.
@pytest.mark.parametrize across empty and multi-item lists.AsyncMock to verify invocation count and arguments.Hypothesis property-based test proving reducer associativity: \[\text{apply}(\text{apply}(s, a), b) = \text{apply}(s, a + b)\]# tests/test_runtime_suite.py
import pytest
from unittest.mock import AsyncMock
from hypothesis import given, strategies as st
from operator import add
# 1. Testing reducers with parametrize
@pytest.mark.parametrize("initial,update,expected", [
([], ["msg1"], ["msg1"]),
(["m1"], ["m2", "m3"], ["m1", "m2", "m3"]),
(["m1"], [], ["m1"]),
])
def test_reducer_merge(initial, update, expected):
assert add(initial, update) == expected
# 2. Testing hooks with AsyncMock
@pytest.mark.asyncio
async def test_hooks_invocation():
mock_hook = AsyncMock()
# Simulate execution hook
await mock_hook("node_fetch", {"data": 123})
mock_hook.assert_awaited_once_with("node_fetch", {"data": 123})# 3. Flaky retry test without real sleeping
@pytest.mark.asyncio
async def test_retry_flaky_node(monkeypatch):
import asyncio
monkeypatch.setattr(asyncio, "sleep", AsyncMock())
attempts = 0
async def flaky_call():
nonlocal attempts
attempts += 1
if attempts < 3:
raise ConnectionError("503 Service Unavailable")
return "recovered"
result = await retry_with_backoff(flaky_call, attempts=4)
assert result == "recovered"
assert attempts == 3
# 4. Hypothesis property test: Associativity
@given(st.lists(st.text()), st.lists(st.text()), st.lists(st.text()))
def test_reducer_associativity(a, b, c):
assert add(add(a, b), c) == add(a, add(b, c))import timeit
import cProfile
import pstats
import tracemalloc
# 1. Micro-benchmarking syntax
timeit.timeit("sum(range(1000))", number=10_000)
# 2. Memory leak detection
tracemalloc.start()
run_agent_pipeline()
current, peak = tracemalloc.get_traced_memory()
print(f"RAM Peak: {peak / 1024 / 1024:.2f} MB")
tracemalloc.stop()
# 3. CPU profiling
cProfile.run("main()", "agent.prof")
stats = pstats.Stats("agent.prof")
stats.sort_stats("cumtime").print_stats(10)from functools import cache, lru_cache
from dataclasses import dataclass
# Unbounded memoization for immutable static schemas
@cache
def load_tool_json_schema(tool_name: str) -> dict:
...
# Bounded LRU cache for document embeddings
@lru_cache(maxsize=1024)
def compute_text_hash(text: str) -> str:
...
# Memory-efficient dataclasses with __slots__
@dataclass(slots=True, frozen=True)
class StreamMessage:
role: str
content: str@cache) on unbounded user inputs cause catastrophic memory leaks in 24/7 agent servers. Always bound with @lru_cache(maxsize=N).@cache require all arguments to be hashable (strings, tuples, frozen dataclasses).| Workload Type | Recommended Architecture | Engineering Rationale |
|---|---|---|
| Many network I/O calls | asyncio |
Low-overhead tasks waiting on remote sockets |
| Blocking legacy libraries | threads / to_thread |
Keeps main asyncio loop responsive |
| CPU-heavy parsing / vector math | multiprocessing |
Bypasses GIL; distributes compute over CPU cores |
| Isolated parallel Python | subinterpreters (3.14) |
Isolated memory and interpreter states in one process |
| Threaded parallel compute | Free-threaded (3.14) | True multi-core execution on modern hardware |
Default Stance: Always start with
asynciofor agent systems. Reach for processes or threads only after measuring a clear CPU bottleneck.
Parallel Tool Calls
Run independent model tool calls concurrently using asyncio.TaskGroup instead of sequentially.
Stream Everything
Stream tokens immediately to reduce user-perceived latency (TTFT) from seconds to milliseconds.
Reuse Connections
Maintain a shared httpx.AsyncClient connection pool to eliminate repeat TLS handshakes.
Budget Tokens
Prune conversational history, truncate repetitive tool outputs, and enforce strict token ceilings.
Cache Pure Lookups
Memoize unchanging prompt templates, tool schemas, and document embeddings.
Monitor Tail Latency
Track p95 and p99 latency metrics; average latency hides stalling agent tasks.
sdist vs Wheel: Source distribution contains raw code; wheel is pre-built, instantly installable without compiling.[project.scripts] turns any Python function into a global shell command when installed.MAJOR: Breaking state or schema changesMINOR: New tools or agent capabilitiesPATCH: Internal bug fixes & performanceimport argparse
import asyncio
from my_agent.runtime import run_agent
def main() -> None:
parser = argparse.ArgumentParser(
prog="my-agent",
description="Autonomous CLI Research Agent"
)
parser.add_argument("query", help="User prompt or task")
parser.add_argument("--model", default="claude-3-5-sonnet")
parser.add_argument("--timeout", type=float, default=60.0)
parser.add_argument("--verbose", "-v", action="store_true")
args = parser.parse_args()
# The single sync-to-async bridge at CLI boundary
result = asyncio.run(
run_agent(args.query, model=args.model, timeout=args.timeout)
)
print(result)
if __name__ == "__main__":
main()argparse: Dependency-free, reliable, built into Python.typer or click for large interactive terminal dashboards.asyncio.run() only once inside main().FROM python:3.13-slim
# Install uv from official Astral binary distribution
COPY --from=ghcr.io/astral-sh/uv:latest /uv /usr/local/bin/uv
WORKDIR /app
# Step 1: Copy lockfile and install dependencies (cached Docker layer!)
COPY pyproject.toml uv.lock ./
RUN uv sync --frozen --no-dev --no-install-project
# Step 2: Copy application source code
COPY src ./src
# Step 3: Complete project installation
RUN uv sync --frozen --no-dev
# Principle of least privilege: run as non-root user
RUN useradd -m -u 1000 appuser
USER appuser
ENV PATH="/app/.venv/bin:$PATH"
ENTRYPOINT ["my-agent"]pyproject.toml and uv.lock first, Docker caches third-party dependencies. Changing Python source code will not trigger re-downloads.--frozen --no-dev: Guarantees repeatable production images without development tools.appuser).Secret Management
Inject API keys via environment variables or secret vaults. Never commit secrets to git or log them in telemetry.
Untrusted Model Output
Never execute unvalidated model outputs. Always parse through strict Pydantic schemas before taking actions.
No eval() or shell=True
Never pass model-generated strings directly into eval(), exec(), or subprocess.Popen(shell=True).
Forbid pickle
Never deserialize untrusted payloads using pickle.loads(). Use JSON or Pydantic models.
Goal
Transform the agent runtime into a deployable, containerized CLI application and verify security compliance.
[project.scripts] entry point and implement an argparse CLI.uv build..whl in a fresh virtual environment and verify the CLI.ghcr.io/astral-sh/uv.uv audit and ensure zero high-severity CVEs exist in the dependency tree.# src/my_agent/cli.py
import argparse
import asyncio
from my_agent.engine import run_pipeline
def main() -> None:
parser = argparse.ArgumentParser(
prog="agent-cli",
description="Autonomous AI Engineering Task Runner"
)
parser.add_argument("task", help="Objective prompt")
parser.add_argument("--retries", type=int, default=3)
args = parser.parse_args()
result = asyncio.run(run_pipeline(args.task, retries=args.retries))
print(f"Task Output:\n{result}")
if __name__ == "__main__":
main()# Packaging & Verification Pipeline:
# 1. Bump version
uv version --bump patch
# 2. Build wheel & sdist
uv build
# 3. Test installation in clean isolated venv
uv venv /tmp/test-venv
uv pip install dist/*.whl --python /tmp/test-venv/bin/python
/tmp/test-venv/bin/agent-cli "Run system healthcheck"
# 4. Security vulnerability scan
uv audit
# 5. Build lean Docker image
docker build -t my-agent:latest .
docker run --rm my-agent:latest "Scan logs"The Challenge
Deliver a small, fully tested, typed, asynchronous agent runtime package that mirrors how LangGraph, Strands, and ADK work internally.
1. Typed State
TypedDict, functional reducers (operator.add), and Annotated state channels.
2. Validated Tools
Automatic JSON Schema extraction from Python functions with Pydantic validation.
3. Async Engine
DAG node execution using graphlib, per-step asyncio.timeout, and cancellation.
4. Extensibility
Async lifecycle hooks (before_tool, after_tool) and Protocol-based dependency injection.
5. Reliability
Retries with exponential backoff and jitter; atomic SQLite state checkpointing.
6. Packaging & Ship
Managed via uv, linted with ruff, strictly checked with mypy, with a CLI entry point.
Suggested Rubric (100 pts): Correctness (30), Typing & Architecture (20), Async & Reliability (20), Tests & CI (15), Packaging & CLI (15).
import asyncio
from graphlib import TopologicalSorter
from operator import add
from typing import Annotated, TypedDict
from pydantic import BaseModel, Field
class AgentState(TypedDict):
task: str
messages: Annotated[list[str], add]
completed: bool
class AgentRuntime:
def __init__(self):
self.nodes = {}
self.deps = {}
self.hooks = []
def node(self, name: str, depends_on: set[str] | None = None):
def decorator(fn):
self.nodes[name] = fn
self.deps[name] = depends_on or set()
return fn
return decorator
async def invoke(self, state: AgentState) -> AgentState:
sorter = TopologicalSorter(self.deps)
current = dict(state)
for name in sorter.static_order():
for h in self.hooks: await h(name, current)
async with asyncio.timeout(5.0):
update = await self.nodes[name](current)
current["messages"] = current.get("messages", []) + update.get("messages", [])
current.update({k: v for k, v in update.items() if k != "messages"})
return current# Executable Application Demonstration:
runtime = AgentRuntime()
@runtime.node("plan")
async def plan_step(state: AgentState):
await asyncio.sleep(0.05)
return {"messages": [f"Plan for: {state['task']}"]}
@runtime.node("execute", depends_on={"plan"})
async def execute_step(state: AgentState):
await asyncio.sleep(0.05)
return {"messages": ["Execution succeeded"], "completed": True}
async def main():
initial: AgentState = {
"task": "Analyze vector embeddings",
"messages": [],
"completed": False
}
final_state = await runtime.invoke(initial)
print("Final State Messages:")
for msg in final_state["messages"]:
print(f" - {msg}")
# asyncio.run(main())# src/my_agent/api/routes.py (Modular APIRouter + Streaming SSE)
from fastapi import APIRouter, Depends, Header, HTTPException, status
from fastapi.responses import StreamingResponse
from pydantic import BaseModel
from my_agent.engine import stream_agent_reasoning
router = APIRouter(prefix="/v1/agent", tags=["Agent"])
class QueryRequest(BaseModel):
prompt: str
temperature: float = 0.7
async def verify_bearer_token(authorization: str = Header(...)):
if not authorization.startswith("Bearer sk-"):
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail="Invalid authentication credentials"
)
return authorization
@router.post("/stream")
async def stream_agent(
req: QueryRequest,
token: str = Depends(verify_bearer_token)
):
"""Industry standard streaming route with Pydantic contract."""
async def sse_event_generator():
async for chunk in stream_agent_reasoning(req.prompt):
yield f"data: {chunk}\n\n"
yield "data: [DONE]\n\n"
return StreamingResponse(
sse_event_generator(),
media_type="text/event-stream"
)# src/my_agent/main.py (Lifespan + Middleware)
from contextlib import asynccontextmanager
from fastapi import FastAPI, Request
from fastapi.middleware.cors import CORSMiddleware
import time
@asynccontextmanager
async def lifespan(app: FastAPI):
# Startup: initialize pooled HTTP clients & database
app.state.db_pool = init_connection_pool()
yield
# Shutdown: cleanly drain connections
await app.state.db_pool.close()
app = FastAPI(title="Agent Service", lifespan=lifespan)
# Global Observability Middleware
@app.middleware("http")
async def add_timing_and_trace(request: Request, call_next):
start = time.perf_counter()
response = await call_next(request)
duration = time.perf_counter() - start
response.headers["X-Response-Time"] = f"{duration:.3f}s"
return response
app.include_router(router)| Python Skill | How It Translates Directly to Frameworks |
|---|---|
TypedDict & Annotated |
Define LangGraph graph state and channel reducers; specify tool schemas |
Pydantic & TypeAdapter |
Parse structured output in Strands; define input/output contracts in ADK |
| Decorators & Generics | Author @tool, @hook, and @node; understand Runtime[Context] and Case[str, str] |
async, gather, TaskGroup |
Execute parallel tool calls, dynamic branch fan-out, and multi-agent coordination |
Async Generators (yield) |
Stream token generation via stream_async() and server-sent events (SSE) |
| Context Managers | Manage Model Context Protocol (MCP) clients, DB transactions, and session state |
| Exceptions & Cancellation | Handle agent stop reasons, human-in-the-loop interrupts, and timeouts safely |
uv, ruff, pytest, CI |
Test, package, containerize, and deploy production-ready autonomous agents |
| Symptom / Error | Root Cause | Idiomatic Resolution |
|---|---|---|
RuntimeWarning: coroutine was never awaited |
Forgot to await a coroutine function |
Add await or schedule with asyncio.create_task() |
RuntimeError: asyncio.run() cannot be called... |
Called asyncio.run() inside an active loop |
Await the coroutine directly |
| Stream freezes and timeouts fire late | Synchronous blocking call inside async function | Offload with await asyncio.to_thread() |
| Shared list grows across independent calls | Mutable default parameter (items=[]) |
Use items=None or field(default_factory=list) |
TypeError: Object is not JSON serializable |
Non-primitive (datetime, set, object) in state | Convert to string/dict or use Pydantic model_dump_json() |
ImportError: cannot import name ... |
Circular dependency between modules | Refactor shared code or guard with if TYPE_CHECKING: |
uv run uses unexpected dependency versions |
Inline script metadata header overrides project | Check # /// script metadata block |
| Type checker passes but runtime crashes | Static annotations are not enforced at runtime | Validate inputs at boundaries using Pydantic |
uvModern Tooling Ecosystem