Skip to content
Merged
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
2 changes: 1 addition & 1 deletion providers/common/ai/README.rst
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,7 @@ PIP package Version required
``apache-airflow`` ``>=3.0.0``
``apache-airflow-providers-common-compat`` ``>=1.15.0``
``apache-airflow-providers-standard`` ``>=1.12.1``
``pydantic-ai-slim`` ``>=1.99.0,<2``
``pydantic-ai-slim`` ``>=2.0.0``
========================================== ==================

Optional cross provider package dependencies
Expand Down
2 changes: 1 addition & 1 deletion providers/common/ai/docs/index.rst
Original file line number Diff line number Diff line change
Expand Up @@ -109,7 +109,7 @@ PIP package Version required
``apache-airflow`` ``>=3.0.0``
``apache-airflow-providers-common-compat`` ``>=1.15.0``
``apache-airflow-providers-standard`` ``>=1.12.1``
``pydantic-ai-slim`` ``>=1.99.0,<2``
``pydantic-ai-slim`` ``>=2.0.0``
========================================== ==================

Optional cross provider package dependencies
Expand Down
8 changes: 8 additions & 0 deletions providers/common/ai/docs/observability.rst
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,14 @@ How it works
names, and finish reason are recorded. Prompt and completion text is never
emitted unless you opt in (see below).

.. note::

On pydantic-ai 2.x the agent-run span reports token usage under
``gen_ai.aggregated_usage.*`` while the per-model-call span keeps
``gen_ai.usage.*``. This avoids double-counting in backends that sum a
parent span and its children. Dashboards or alerts that read run-level token
usage from ``gen_ai.usage.*`` should switch to ``gen_ai.aggregated_usage.*``.

Enabling it
-----------

Expand Down
10 changes: 4 additions & 6 deletions providers/common/ai/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -69,9 +69,8 @@ dependencies = [
"apache-airflow>=3.0.0",
"apache-airflow-providers-common-compat>=1.15.0",
"apache-airflow-providers-standard>=1.12.1",
# Capped to 1.x: pydantic-ai 2.x changed the agent instrumentation API and breaks the provider.
# Remove the cap after migrating; tracked at https://github.com/apache/airflow/issues/69122
"pydantic-ai-slim>=1.99.0,<2",
# Requires the pydantic-ai 2.x agent/instrumentation API (see #69122).
"pydantic-ai-slim>=2.0.0",
]

# The optional dependencies should be modified in place in the generated file
Expand All @@ -87,9 +86,8 @@ dependencies = [
# Enables AgentOperator(code_mode=True). Monty is pre-1.0; pinned here as an
# opt-in extra so its churn never breaks the base provider install.
"code-mode" = ["pydantic-ai-harness[codemode]>=0.3.0"]
# Agent Skills (agentskills.io) support. pydantic-ai-skills provides the toolset
# (pulls in pydantic-ai-slim>=1.74 transitively; the provider base floor stays
# 1.71); the git provider supplies GitHook + GitPython for cloning GitSkills with
# Agent Skills (agentskills.io) support. pydantic-ai-skills provides the toolset;
# the git provider supplies GitHook + GitPython for cloning GitSkills with
# credentials from a `git` connection. Native progressive disclosure is tracked
# upstream in pydantic/pydantic-ai#5230; revisit this extra once that lands.
"skills" = [
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,11 @@

OutputT = TypeVar("OutputT")

# Sentinel distinguishing "caller did not pass ``instrument``" from an explicit
# ``instrument=None`` / ``instrument=False`` (which mean "do not instrument, and
# do not auto-enable it either").
_UNSET: Any = object()

if TYPE_CHECKING:
from pydantic_ai.models import KnownModelName, Model

Expand Down Expand Up @@ -196,10 +201,10 @@ def _get_conn_if_model_configured(self) -> Model | None:
@overload
def create_agent(
self, output_type: type[OutputT], *, instructions: str, **agent_kwargs
) -> Agent[None, OutputT]: ...
) -> Agent[object, OutputT]: ...

@overload
def create_agent(self, *, instructions: str, **agent_kwargs) -> Agent[None, str]: ...
def create_agent(self, *, instructions: str, **agent_kwargs) -> Agent[object, str]: ...

@overload
def create_agent(
Expand All @@ -209,7 +214,7 @@ def create_agent(
spec_file: str | Path,
instructions: str | None = ...,
**agent_kwargs,
) -> Agent[None, OutputT]: ...
) -> Agent[object, OutputT]: ...

@overload
def create_agent(
Expand All @@ -218,7 +223,7 @@ def create_agent(
spec_file: str | Path,
instructions: str | None = ...,
**agent_kwargs,
) -> Agent[None, str]: ...
) -> Agent[object, str]: ...

def create_agent(
self,
Expand All @@ -227,7 +232,7 @@ def create_agent(
instructions: str | None = None,
spec_file: str | Path | None = None,
**agent_kwargs,
) -> Agent[None, Any]:
) -> Agent[object, Any]:
"""
Create a pydantic-ai Agent configured with this hook's model.

Expand All @@ -247,6 +252,14 @@ def create_agent(
the spec file's ``model`` is used.
:param agent_kwargs: Additional keyword arguments passed to the Agent constructor.
"""
# ``instrument`` is no longer an ``Agent()`` / ``Agent.from_file()``
# constructor argument in pydantic-ai 2.x; it is configured through the
# ``agent.instrument`` property (which is unchanged across the 2.x line).
# Pop any caller-supplied value out of the constructor kwargs and apply
# it after construction so a caller that passes its own ``instrument``
# still wins over the provider's auto-instrumentation.
caller_instrument = agent_kwargs.pop("instrument", _UNSET)

if spec_file is not None:
from_file_kwargs = dict(agent_kwargs)
model = self._get_conn_if_model_configured()
Expand All @@ -264,13 +277,10 @@ def create_agent(
if instructions is None:
raise ValueError("instructions is required when spec_file is not provided.")
agent = Agent(self.get_conn(), output_type=output_type, instructions=instructions, **agent_kwargs)
if "instrument" not in agent_kwargs:
# Set the public ``agent.instrument`` surface rather than the
# ``Agent(instrument=...)`` constructor kwarg, which is deprecated in
# current pydantic-ai. Assigning ``agent.instrument`` works across the
# provider's ``pydantic-ai-slim>=1.71`` floor (a plain instance
# attribute on older versions, a property on newer ones). A caller
# that passed its own ``instrument`` wins.

if caller_instrument is not _UNSET:
agent.instrument = caller_instrument
else:
settings = genai_instrumentation_settings()
if settings is not None:
agent.instrument = settings
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -40,10 +40,13 @@

SECTION = "common.ai"

# OTel GenAI semantic-convention attribute set. Pinned so the emitted span and
# ``gen_ai.*`` attribute names stay stable regardless of the pydantic-ai default
# (which tracks the latest, still-evolving revision).
_SEMCONV_VERSION: Literal[4] = 4
# OTel GenAI semantic-convention format version. Pinned so a change in the
# pydantic-ai default does not silently shift the emitted span/attribute format
# between provider releases. Version 5 is the current default in pydantic-ai 2.x;
# formats 2-4 still work but are deprecated. Note: independent of this version,
# pydantic-ai 2.x reports agent-run token usage under ``gen_ai.aggregated_usage.*``
# (model-request spans keep ``gen_ai.usage.*``) -- see docs/observability.rst.
_SEMCONV_VERSION: Literal[5] = 5


def _otel_export_enabled() -> bool:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -301,7 +301,7 @@ def llm_hook(self) -> PydanticAIHook:
}
return PydanticAIHook.get_hook(self.llm_conn_id, hook_params=hook_params)

def _build_agent(self) -> Agent[None, Any]:
def _build_agent(self) -> Agent[object, Any]:
"""Build and return a pydantic-ai Agent from the operator's config."""
extra_kwargs = dict(self.agent_params)
if self.toolsets:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -163,7 +163,7 @@ def execute(self, context: Context) -> Any:
f"str prompt, or disable require_approval."
)

agent: Agent[None, Any] = self.llm_hook.create_agent(
agent: Agent[object, Any] = self.llm_hook.create_agent(
output_type=self.output_type, instructions=self.system_prompt, **self.agent_params
)
result = agent.run_sync(self.prompt, usage_limits=self.usage_limits)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -129,7 +129,7 @@ def execute(self, context: Context) -> Any:
self.sample_rows,
)
self.log.debug("Resolved file analysis paths: %s", request.resolved_paths)
agent: Agent[None, Any] = self.llm_hook.create_agent(
agent: Agent[object, Any] = self.llm_hook.create_agent(
output_type=self.output_type,
instructions=self._build_system_prompt(),
**self.agent_params,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -404,12 +404,23 @@ def test_no_instrument_when_settings_none(self, mock_settings):
@patch("airflow.providers.common.ai.hooks.pydantic_ai.infer_model", autospec=True)
def test_caller_instrument_short_circuits(self, mock_infer_model, mock_settings, mock_agent_cls):
"""A caller that passes its own ``instrument`` wins; we don't override it."""
mock_infer_model.return_value = MagicMock(spec=Model)
mock_model = MagicMock(spec=Model)
mock_infer_model.return_value = mock_model
hook = self._hook()
conn = Connection(conn_id="test_conn", conn_type="pydanticai")
with patch.object(hook, "get_connection", return_value=conn):
hook.create_agent(instructions="hi", instrument=False)
agent = hook.create_agent(instructions="hi", instrument=False)

# ``instrument`` is not an Agent() constructor kwarg in pydantic-ai 2.x:
# it must be stripped from the constructor call and applied through the
# ``agent.instrument`` property instead, and the provider's own
# auto-instrumentation must not override the caller's value.
mock_agent_cls.assert_called_once_with(
mock_model,
output_type=str,
instructions="hi",
)
assert agent.instrument is False
mock_settings.assert_not_called()


Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,7 @@ def test_returns_settings_when_enabled_with_provider(self):
settings = observability.genai_instrumentation_settings()

assert settings is not None
assert settings.version == observability._SEMCONV_VERSION == 4
assert settings.version == observability._SEMCONV_VERSION == 5
# Content is never emitted unless explicitly opted in.
assert settings.include_content is False
assert settings.include_binary_content is False
Expand Down Expand Up @@ -123,8 +123,13 @@ def test_spans_emitted_and_nested_without_content_by_default(self):
genai, parent_trace_id, attrs_blob = self._run(capture=False)

assert genai, "expected gen_ai spans to be emitted"
# Token usage is captured even with content off.
# Token usage is captured even with content off. In pydantic-ai 2.x the
# model-request span keeps ``gen_ai.usage.*`` while the agent-run span
# reports ``gen_ai.aggregated_usage.*`` (avoids double-counting when a
# backend sums parent and child spans); assert both so a change in that
# split is caught here rather than silently shifting users' telemetry.
assert any("gen_ai.usage.input_tokens" in (s.attributes or {}) for s in genai)
assert any("gen_ai.aggregated_usage.input_tokens" in (s.attributes or {}) for s in genai)
# Parenting is implicit: agent spans share the task span's trace_id.
assert all(s.context.trace_id == parent_trace_id for s in genai)
# The prompt text must not leak when content capture is off.
Expand Down
14 changes: 7 additions & 7 deletions uv.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading