Repository navigation
Conversation
| for run in runs: | ||
| key = run.created_dag_version_id if run.bundle_version and run.created_dag_version_id else None | ||
| representatives.setdefault(key, run) | ||
| for run in representatives.values(): |
There was a problem hiding this comment.
Agreed, thanks. The group is now resolved per run: _get_dags_for_named_runs maps each existing named run to the version it resolves to (one deserialization per distinct version), _get_group_task_ids_by_run keeps the task ids per run, and _get_group_tasks filters with or_(and_(TI.run_id == run_id, TI.task_id.in_(ids))), so a task that is part of the group in one version is never attributed to a run of another. The latest version only answers when no named run exists. The count endpoint had the same leak one layer up: it reduced the group instances to (task_id, map_index) pairs and matched them across every named run, so it now matches (run_id, task_id, map_index). New tests cover the prefix_group_id=False case you describe on both endpoints (task2 outside the group in v1 and inside it in v2: 3 instances for the two runs, not 4).
The same per-run resolution now also gives the provider side (#73085) the answer it needs: task-instances/states returns 404 for a task_ids entry that no named run's version defines and no named run has an instance of. A task the version defines counts before its instance exists (the reparse window the scheduler's _verify_integrity_if_dag_changed leaves), and an instance counts on its own (a task removed in a newer version), so existing callers that ask about their own task are unaffected. Runs that do not exist are skipped, following the point on #67832 that nothing is known before a run exists. count is left alone, since the task runner uses it for mapped-task counts.
| for run in runs: | ||
| key = run.created_dag_version_id if run.bundle_version and run.created_dag_version_id else None | ||
| representatives.setdefault(key, run) | ||
| for run in representatives.values(): |
There was a problem hiding this comment.
Agreed, thanks. The group is now resolved per run: _get_dags_for_named_runs maps each existing named run to the version it resolves to (one deserialization per distinct version), _get_group_task_ids_by_run keeps the task ids per run, and _get_group_tasks filters with or_(and_(TI.run_id == run_id, TI.task_id.in_(ids))), so a task that is part of the group in one version is never attributed to a run of another. The latest version only answers when no named run exists. The count endpoint had the same leak one layer up: it reduced the group instances to (task_id, map_index) pairs and matched them across every named run, so it now matches (run_id, task_id, map_index). New tests cover the prefix_group_id=False case you describe on both endpoints (task2 outside the group in v1 and inside it in v2: 3 instances for the two runs, not 4).
The same per-run resolution now also gives the provider side (#73085) the answer it needs: task-instances/states returns 404 for a task_ids entry that no named run's version defines and no named run has an instance of. A task the version defines counts before its instance exists (the reparse window the scheduler's _verify_integrity_if_dag_changed leaves), and an instance counts on its own (a task removed in a newer version), so existing callers that ask about their own task are unaffected. Runs that do not exist are skipped, following the point on #67832 that nothing is known before a run exists. count is left alone, since the task runner uses it for mapped-task counts.
c7874df to
22642d6
Compare
…sions Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
22642d6 to
294846d
Compare
|
Rebased onto main. The only conflict was the import block that #73554 removed; @ashb this follows your suggestion on #67832: instead of a separate existence endpoint, @Vamsi-klu the per-run resolution you asked for is in; the details are in your thread. |
_get_group_tasks, which backs thetask_group_idfilter ofGET /execution/task-instances/countandGET /execution/task-instances/states, resolves the task group against the latest Dag version, even when the request names the runs it is asking about. For a run of a versioned bundle that is the wrong version whenever the group changed after the run was created: a group renamed or removed since answers 404 although the run has its task instances, and a group that only exists in the newer version is reported as part of the run although it is not.ExternalTaskSensorandWorkflowTriggerpoll these endpoints withlogical_dates/run_ids, and the existence check for #72514 (#73085) needs the answer to be about the named runs.Task groups are resolved per run. Each named run that exists is mapped to the Dag version it resolves to through
DBDagBag.get_dag_for_run(the version a run of a versioned bundle was created from, the latest version for every other run, which is also how the scheduler treats those runs when it re-verifies them), with one deserialization per distinct version. Each run then contributes only the tasks its own version places in the group, so a task that is part of the group in one version is never attributed to a run of another; this matters withprefix_group_id=False, where task ids do not carry the group name. The count endpoint had the same leak one layer up: it reduced the group instances to(task_id, map_index)pairs and matched them across every named run, so it now matches(run_id, task_id, map_index). Without a named run, or while none of the named runs exists yet, the latest version answers as before.Unknown tasks are reported for named runs.
GET /execution/task-instances/statesnow answers 404 (reason: not_found) for atask_idsentry that no named run's Dag version defines and no named run has an instance of. A task the version defines exists for the run before its instance is created (the scheduler adds newly parsed tasks to an unfinished run of an unversioned bundle on a later pass), and an existing instance counts on its own (a task removed in a newer version is still found for a run that has it), so callers that ask about their own task are unaffected. Runs that do not exist yet are skipped, and a request that names no run, or only runs that do not exist, is not validated: nothing is known about a run before it exists, the point made on #67832 against a separate existence endpoint. This is the answer #73085 relies on, in the request the sensor already makes.countis unchanged, since the task runner uses it for mapped-task counts. No execution API version change: request and response schemas are the same, and the new 404 uses the existing payload, as #63355 did for its 409.Changes
airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py:_get_dags_for_named_runs,_task_ids_in_groupand_get_group_task_ids_by_run(per-run group resolution);_get_group_tasksfilters per run; the count endpoint matches(run_id, task_id, map_index);_raise_if_tasks_unknown_to_runsin the states endpointairflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py:test_get_count_task_group_resolved_against_run_versionandtest_get_task_states_task_group_resolved_against_run_version(a pinned run keeps its group after a rename);test_get_count_task_group_resolved_per_runandtest_get_task_states_task_group_resolved_per_run(prefix_group_id=False, a task outside the group in v1 and inside it in v2: 3 instances for the two runs, not 4);test_get_task_states_unknown_task_for_named_run(404 with one or several missing tasks; no validation for a run that does not exist or without a named run);test_get_task_states_task_known_to_run_version_or_by_instance(a task added in the latest version is known before its instance exists, a removed one through its instance). All fail on main.Testing
airflow-core:tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py(TestGetCountandTestGetTaskStates, 46 tests)related: #72514. Needed by #73085.
Was generative AI tooling used to co-author this PR?
Generated-by: Claude Code (Claude Fable 5.1) following the guidelines. I reviewed and understand all changes; the tests were run locally as listed above.
🤖 Generated with Claude Code