This is an automated email from the ASF dual-hosted git repository.
pierrejeambrun pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 14960cf4624 Add AIP-85 to the AIP progress tracker registry (#74304)
14960cf4624 is described below
commit 14960cf4624a030310062e96af2e751b64ebc1e9
Author: Pierre Jeambrun <[email protected]>
AuthorDate: Tue Oct 6 11:08:34 2026 +0200
Add AIP-85 to the AIP progress tracker registry (#74304)
* Add AIP-85 to the AIP progress tracker registry
Tracking runs scoped to AIP-85 had no registry entry to work from, so
the agent fell back to an unconstrained search instead of AIP-85's
Confluence page and known codebase paths, and overran its context
budget investigating every other tracked AIP along the way.
* Respect aip_numbers when building the skills-based tracker's prompt
The Dag-decorated function runs once at parse time, before any run's
aip_numbers conf exists, so building the agent's prompt directly in
the Dag body silently ignored the conf and always listed every
registered AIP -- regardless of what a run actually asked for. This
made runs scoped to a single AIP still investigate every tracked AIP,
which is what drove AIP-85 runs over their context budget. Move the
prompt construction into a task, which runs per-DagRun and can read
the resolved conf, matching how the deterministic tracker Dag already
scopes its own AIP list.
---
.../example_dags/example_aip_progress_tracker.py | 64 +++++++++++++++++-----
1 file changed, 50 insertions(+), 14 deletions(-)
diff --git
a/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_aip_progress_tracker.py
b/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_aip_progress_tracker.py
index 5a1e82de0b4..191935e42aa 100644
---
a/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_aip_progress_tracker.py
+++
b/providers/common/ai/src/airflow/providers/common/ai/example_dags/example_aip_progress_tracker.py
@@ -177,6 +177,33 @@ AIP_REGISTRY: dict[int, dict] = {
"java-sdk",
],
},
+ 85: {
+ "page_id": "315494137",
+ "topic": "Dag Importer (native Dag parsing for Lang-SDK languages)",
+ "search_terms": [
+ "Dag importer",
+ "native Dag",
+ "NodeDagImporter",
+ "JavaDagImporter",
+ "lang_sdk_processor",
+ "CoordinatorDagImporter",
+ "importer registry",
+ "DagDef",
+ ],
+ "codebase_paths": [
+ "airflow-core/src/airflow/dag_processing/lang_sdk_processor.py",
+ "airflow-core/src/airflow/dag_processing/importer_routing.py",
+ "task-sdk/src/airflow/sdk/coordinators/_dag_importer.py",
+ "task-sdk/src/airflow/sdk/coordinators/node",
+ "task-sdk/src/airflow/sdk/coordinators/java",
+ "task-sdk/src/airflow/sdk/importers/base.py",
+ "airflow-core/adr/lang-sdk/0010-native-dag-processing.md",
+ "kubernetes-tests/lang_sdk",
+ "ts-sdk/src/sdk",
+ "go-sdk/airflow",
+ "java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk",
+ ],
+ },
}
# [END aip_registry]
@@ -1097,31 +1124,39 @@ def example_aip_progress_tracker_skills():
tools=[fetch_confluence_page, search_github_prs, get_repo_file_tree],
)
- aip_info = "\n".join(
- f"- AIP-{num}: {info['topic']} (page_id={info['page_id']}, "
- f"paths: {', '.join(info['codebase_paths'][:3])})"
- for num, info in AIP_REGISTRY.items()
- )
+ # Built in a task, not at parse time: the Dag function runs once at parse
time, before a
+ # run's `aip_numbers` conf exists, so filtering the registry needs to
happen at task
+ # execution. The result reaches AgentOperator through its templated
`prompt` field.
+ @task
+ def build_agent_prompt(params: dict) -> str:
+ requested = {int(n.strip()) for n in params["aip_numbers"].split(",")
if n.strip()}
+ scoped_registry = {num: info for num, info in AIP_REGISTRY.items() if
num in requested}
+ aip_info = "\n".join(
+ f"- AIP-{num}: {info['topic']} (page_id={info['page_id']}, "
+ f"paths: {', '.join(info['codebase_paths'][:3])})"
+ for num, info in scoped_registry.items()
+ )
+ return (
+ "Track the implementation progress of these AIPs and produce a "
+ "cross-AIP progress report:\n\n"
+ f"{aip_info}\n\n"
+ "Use the aip-tracker skill for detailed instructions on how to "
+ "gather evidence and structure your assessment."
+ )
- prompt = (
- "Track the implementation progress of these AIPs and produce a "
- "cross-AIP progress report:\n\n"
- f"{aip_info}\n\n"
- "Use the aip-tracker skill for detailed instructions on how to "
- "gather evidence and structure your assessment."
- )
+ agent_prompt = build_agent_prompt()
# [START aip_tracker_skills_operator]
report = AgentOperator(
task_id="track_aip_progress",
llm_conn_id=LLM_CONN_ID,
system_prompt=AGENT_SYSTEM_PROMPT,
- prompt=prompt,
+ prompt="{{ ti.xcom_pull(task_ids='build_agent_prompt') }}",
toolsets=[
AgentSkillsToolset(sources=[SKILLS_DIR]),
aip_toolset,
],
- agent_params={"model_settings": {"temperature": 0}},
+ agent_params={"model_settings": {"temperature": 0, "max_tokens":
16_000}},
usage_limits=UsageLimits(
request_limit=30,
input_tokens_limit=200_000,
@@ -1129,6 +1164,7 @@ def example_aip_progress_tracker_skills():
),
)
# [END aip_tracker_skills_operator]
+ agent_prompt >> report
# [START aip_tracker_skills_hitl]
ApprovalOperator(