Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
89 changes: 89 additions & 0 deletions examples/otel_tracing.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,89 @@
#!/usr/bin/env python3
"""Example: OpenTelemetry tracing with the Claude Agent SDK.

This example shows how to wire up distributed tracing so every SDK
call -- session start, message, tool invocation -- appears as a span
in your observability backend (Jaeger, Zipkin, OTLP-compatible, ...).

Prerequisites
-------------
Install the SDK with the ``[otel]`` extra and a span exporter::

pip install claude-agent-sdk[otel] \
opentelemetry-sdk \
opentelemetry-exporter-otlp-proto-grpc

Then run a local Jaeger instance (the all-in-one Docker image is the
fastest way to get a collector + UI)::

docker run -d --name jaeger \
-p 16686:16686 \
-p 4317:4317 \
jaegertracing/all-in-one:latest

Finally, run this script::

python examples/otel_tracing.py

Open http://localhost:16686 to browse the traces in the Jaeger UI.
"""

import anyio

from claude_agent_sdk import (
AssistantMessage,
ClaudeAgentOptions,
ResultMessage,
TextBlock,
enable_tracing,
query,
)


def setup_otel() -> None:
"""Configure an OpenTelemetry TracerProvider with an OTLP exporter."""
from opentelemetry import trace
from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter
from opentelemetry.sdk.resources import Resource
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor

resource = Resource.create({"service.name": "claude-agent-sdk-example"})
provider = TracerProvider(resource=resource)
exporter = OTLPSpanExporter(endpoint="http://localhost:4317", insecure=True)
provider.add_span_processor(BatchSpanProcessor(exporter))
trace.set_tracer_provider(provider)


async def traced_query() -> None:
"""Run a simple query with tracing enabled."""
print("Running a traced query...")

options = ClaudeAgentOptions(max_turns=1)
async for message in query(
prompt="What is the square root of 144? Answer in one sentence.",
options=options,
):
if isinstance(message, AssistantMessage):
for block in message.content:
if isinstance(block, TextBlock):
print(f"Claude: {block.text}")
elif isinstance(message, ResultMessage):
print(f"Turns: {message.num_turns}, Cost: ${message.total_cost_usd:.4f}")

print("\nDone. Check Jaeger UI at http://localhost:16686")


def main() -> None:
# 1. Configure the OTel SDK (provider, exporter, processor).
setup_otel()

# 2. Tell the Claude Agent SDK to start emitting spans.
enable_tracing()

# 3. Run the traced query.
anyio.run(traced_query)


if __name__ == "__main__":
main()
3 changes: 3 additions & 0 deletions src/claude_agent_sdk/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,7 @@
list_subagents,
list_subagents_from_store,
)
from ._internal.tracing import disable_tracing, enable_tracing
from ._internal.transport import Transport
from ._version import __version__
from .client import ClaudeSDKClient
Expand Down Expand Up @@ -528,6 +529,8 @@ async def call_tool(name: str, arguments: dict[str, Any]) -> Any:
__all__ = [
# Main exports
"query",
"enable_tracing",
"disable_tracing",
"__version__",
# Transport
"Transport",
Expand Down
42 changes: 42 additions & 0 deletions src/claude_agent_sdk/_internal/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
HookEvent,
HookMatcher,
Message,
ResultMessage,
_warn_if_can_use_tool_shadowed,
)
from .message_parser import parse_message
Expand All @@ -22,6 +23,7 @@
materialize_resume_session,
)
from .session_store_validation import validate_session_store_options
from .tracing import record_span_event, set_span_attributes, start_span
from .transport import Transport
from .transport.subprocess_cli import SubprocessCLITransport

Expand Down Expand Up @@ -95,6 +97,27 @@ async def _process_query_inner(
options: ClaudeAgentOptions,
transport: Transport | None,
materialized: MaterializedResume | None,
) -> AsyncGenerator[Message, None]:
# Build span attributes for the query-level tracing span.
span_attrs: dict[str, Any] = {}
if options.model:
span_attrs["claude_agent_sdk.model"] = options.model
if options.max_turns is not None:
span_attrs["claude_agent_sdk.max_turns"] = options.max_turns

with start_span("claude_agent_sdk.query", attributes=span_attrs) as query_span:
async for msg in self._process_query_inner_traced(
prompt, options, transport, materialized, query_span
):
yield msg

async def _process_query_inner_traced(
self,
prompt: str | AsyncIterable[dict[str, Any]],
options: ClaudeAgentOptions,
transport: Transport | None,
materialized: MaterializedResume | None,
query_span: Any,
) -> AsyncGenerator[Message, None]:
# Validate and configure permission settings (matching TypeScript SDK logic)
configured_options = options
Expand Down Expand Up @@ -227,6 +250,25 @@ async def _on_mirror_error(key: Any, error: str) -> None:
async for data in query.receive_messages():
message = parse_message(data)
if message is not None:
record_span_event(
query_span,
f"message.{type(message).__name__}",
)
if isinstance(message, ResultMessage):
attrs: dict[str, Any] = {}
if message.num_turns is not None:
attrs["claude_agent_sdk.num_turns"] = message.num_turns
if message.is_error is not None:
attrs["claude_agent_sdk.is_error"] = message.is_error
if message.duration_ms is not None:
attrs["claude_agent_sdk.duration_ms"] = message.duration_ms
if message.session_id:
attrs["claude_agent_sdk.session_id"] = message.session_id
if message.total_cost_usd is not None:
attrs["claude_agent_sdk.total_cost_usd"] = (
message.total_cost_usd
)
set_span_attributes(query_span, attrs)
yield message

finally:
Expand Down
68 changes: 44 additions & 24 deletions src/claude_agent_sdk/_internal/query.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@
ToolPermissionContext,
)
from ._task_compat import TaskHandle, spawn_detached
from .tracing import start_span
from .transport import Transport

if TYPE_CHECKING:
Expand Down Expand Up @@ -428,33 +429,38 @@ async def _handle_control_request(self, request: SDKControlRequest) -> None:

if subtype == "can_use_tool":
permission_request: SDKControlPermissionRequest = request_data # type: ignore[assignment]
tool_name = permission_request["tool_name"]
original_input = permission_request["input"]
# Handle tool permission request
if not self.can_use_tool:
raise Exception("canUseTool callback is not provided")

context = ToolPermissionContext(
signal=None, # TODO: Add abort signal support
suggestions=[
PermissionUpdate.from_dict(s)
for s in (
permission_request.get("permission_suggestions") or []
)
],
tool_use_id=permission_request.get("tool_use_id"),
agent_id=permission_request.get("agent_id"),
blocked_path=permission_request.get("blocked_path"),
decision_reason=permission_request.get("decision_reason"),
title=permission_request.get("title"),
display_name=permission_request.get("display_name"),
description=permission_request.get("description"),
)
with start_span(
"claude_agent_sdk.tool_permission",
attributes={"tool.name": tool_name},
):
context = ToolPermissionContext(
signal=None, # TODO: Add abort signal support
suggestions=[
PermissionUpdate.from_dict(s)
for s in (
permission_request.get("permission_suggestions") or []
)
],
tool_use_id=permission_request.get("tool_use_id"),
agent_id=permission_request.get("agent_id"),
blocked_path=permission_request.get("blocked_path"),
decision_reason=permission_request.get("decision_reason"),
title=permission_request.get("title"),
display_name=permission_request.get("display_name"),
description=permission_request.get("description"),
)

response = await self.can_use_tool(
permission_request["tool_name"],
permission_request["input"],
context,
)
response = await self.can_use_tool(
tool_name,
permission_request["input"],
context,
)

# Convert PermissionResult to expected dict format
if isinstance(response, PermissionResultAllow):
Expand Down Expand Up @@ -507,9 +513,23 @@ async def _handle_control_request(self, request: SDKControlRequest) -> None:
# Type narrowing - we've verified these are not None above
assert isinstance(server_name, str)
assert isinstance(mcp_message, dict)
mcp_response = await self._handle_sdk_mcp_request(
server_name, mcp_message
)
mcp_method = mcp_message.get("method", "")
mcp_tool_name = ""
if mcp_method == "tools/call":
params = mcp_message.get("params", {})
if isinstance(params, dict):
mcp_tool_name = params.get("name", "")
with start_span(
"claude_agent_sdk.tool_call",
attributes={
"mcp.server": server_name,
"mcp.method": mcp_method,
"tool.name": mcp_tool_name,
},
):
mcp_response = await self._handle_sdk_mcp_request(
server_name, mcp_message
)
# Wrap the MCP response as expected by the control protocol
response_data = {"mcp_response": mcp_response}

Expand Down
Loading