Uh oh!
There was an error while loading. Please reload this page.
Add endpoint to watch dag run until finish - #51920
Conversation
potiuk
commented
Jun 19, 2025
FYI. Not sure if it helps but OpenAPI 3.2.0 will have support for |
231597b to
007a051Compareuranusjr
commented
Jun 25, 2025
I added some tests for the endpoint, but couldn’t figure out how to test the looping part. Hopefully this is good enough… |
pierrejeambrun
left a comment
There was a problem hiding this comment.
Nice, just a few suggestions / nits.
Indeed an additional test case for running / success state would be great.
Uh oh!
There was an error while loading. Please reload this page.
Uh oh!
There was an error while loading. Please reload this page.
Uh oh!
There was an error while loading. Please reload this page.
Uh oh!
There was an error while loading. Please reload this page.
007a051 to
152bd6dCompare
pierrejeambrun
left a comment
There was a problem hiding this comment.
LGTM, thanks.
Just a few nits, but nothing blocking
Uh oh!
There was an error while loading. Please reload this page.
Uh oh!
There was an error while loading. Please reload this page.
Uh oh!
There was an error while loading. Please reload this page.
152bd6d to
34badc1Compare96812e6 to
422cd94CompareUh oh!
There was an error while loading. Please reload this page.
kaxil
left a comment
There was a problem hiding this comment.
Some simple docs would be good. Fine as a follow-up too as long as we have a GH issue to not forget.
@vikramkoka could review that too
422cd94 to
58d7fa5CompareTracking doc addition #53067 |
(from https://github.com/apache/airflow/tree/python-client/3.1.0rc1) ## New Features: - Add `map_index` filter to TaskInstance API queries ([#55614](apache/airflow#55614)) - Add `has_import_errors` filter to Core API GET /dags endpoint ([#54563](apache/airflow#54563)) - Add `dag_version` filter to get_dag_runs endpoint ([#54882](apache/airflow#54882)) - Implement pattern search for event log endpoint ([#55114](apache/airflow#55114)) - Add asset-based filtering support to DAG API endpoint ([#54263](apache/airflow#54263)) - Add Greater Than and Less Than range filters to DagRuns and Task Instance list ([#54302](apache/airflow#54302)) - Add `try_number` as filter to task instances ([#54695](apache/airflow#54695)) - Add filters to Browse XComs endpoint ([#54049](apache/airflow#54049)) - Add Filtering by DAG Bundle Name and Version to API routes ([#54004](apache/airflow#54004)) - Add search filter for DAG runs by triggering user name ([#53652](apache/airflow#53652)) - Enable multi sorting (AIP-84) ([#53408](apache/airflow#53408)) - Add `run_on_latest_version` support for backfill and clear operations ([#52177](apache/airflow#52177)) - Add `run_id_pattern` search for Dag Run API ([#52437](apache/airflow#52437)) - Add tracking of triggering user to Dag runs ([#51738](apache/airflow#51738)) - Expose DAG parsing duration in the API ([#54752](apache/airflow#54752)) ## New API Endpoints: - Add Human-in-the-Loop (HITL) endpoints for approval workflows ([#52868](apache/airflow#52868), [#53373](apache/airflow#53373), [#53376](apache/airflow#53376), [#53885](apache/airflow#53885), [#53923](apache/airflow#53923), [#54308](apache/airflow#54308), [#54310](apache/airflow#54310), [#54723](apache/airflow#54723), [#54773](apache/airflow#54773), [#55019](apache/airflow#55019), [#55463](apache/airflow#55463), [#55525](apache/airflow#55525), [#55535](apache/airflow#55535), [#55603](apache/airflow#55603), [#55776](apache/airflow#55776)) - Add endpoint to watch dag run until finish ([#51920](apache/airflow#51920)) - Add TI bulk actions endpoint ([#50443](apache/airflow#50443)) - Add Keycloak Refresh Token Endpoint ([#51657](apache/airflow#51657)) ## Deprecations: - Mark `DagDetailsResponse.concurrency` as deprecated ([#55150](apache/airflow#55150)) ## Bug Fixes: - Fix dag import error modal pagination ([#55719](apache/airflow#55719))
(from https://github.com/apache/airflow/tree/python-client/3.1.0rc1) ## New Features: - Add `map_index` filter to TaskInstance API queries ([#55614](apache/airflow#55614)) - Add `has_import_errors` filter to Core API GET /dags endpoint ([#54563](apache/airflow#54563)) - Add `dag_version` filter to get_dag_runs endpoint ([#54882](apache/airflow#54882)) - Implement pattern search for event log endpoint ([#55114](apache/airflow#55114)) - Add asset-based filtering support to DAG API endpoint ([#54263](apache/airflow#54263)) - Add Greater Than and Less Than range filters to DagRuns and Task Instance list ([#54302](apache/airflow#54302)) - Add `try_number` as filter to task instances ([#54695](apache/airflow#54695)) - Add filters to Browse XComs endpoint ([#54049](apache/airflow#54049)) - Add Filtering by DAG Bundle Name and Version to API routes ([#54004](apache/airflow#54004)) - Add search filter for DAG runs by triggering user name ([#53652](apache/airflow#53652)) - Enable multi sorting (AIP-84) ([#53408](apache/airflow#53408)) - Add `run_on_latest_version` support for backfill and clear operations ([#52177](apache/airflow#52177)) - Add `run_id_pattern` search for Dag Run API ([#52437](apache/airflow#52437)) - Add tracking of triggering user to Dag runs ([#51738](apache/airflow#51738)) - Expose DAG parsing duration in the API ([#54752](apache/airflow#54752)) ## New API Endpoints: - Add Human-in-the-Loop (HITL) endpoints for approval workflows ([#52868](apache/airflow#52868), [#53373](apache/airflow#53373), [#53376](apache/airflow#53376), [#53885](apache/airflow#53885), [#53923](apache/airflow#53923), [#54308](apache/airflow#54308), [#54310](apache/airflow#54310), [#54723](apache/airflow#54723), [#54773](apache/airflow#54773), [#55019](apache/airflow#55019), [#55463](apache/airflow#55463), [#55525](apache/airflow#55525), [#55535](apache/airflow#55535), [#55603](apache/airflow#55603), [#55776](apache/airflow#55776)) - Add endpoint to watch dag run until finish ([#51920](apache/airflow#51920)) - Add TI bulk actions endpoint ([#50443](apache/airflow#50443)) - Add Keycloak Refresh Token Endpoint ([#51657](apache/airflow#51657)) ## Deprecations: - Mark `DagDetailsResponse.concurrency` as deprecated ([#55150](apache/airflow#55150)) ## Bug Fixes: - Fix dag import error modal pagination ([#55719](apache/airflow#55719))
When _configure_async_session() was extracted from configure_orm() in PR apache#51920, the corresponding cleanup in dispose_orm() was not updated. This left async_engine connections abandoned on process exit, gunicorn worker restarts, and atexit -- gradually exhausting PostgreSQL's max_connections. Dispose async_engine.sync_engine (the synchronous path, matching the existing clean_in_fork pattern) and clear both async_engine and AsyncSession references.
When _configure_async_session() was extracted from configure_orm() in PR apache#51920, the corresponding cleanup in dispose_orm() was not updated. This left async_engine connections abandoned on process exit, gunicorn worker restarts, and atexit -- gradually exhausting PostgreSQL's max_connections. Dispose async_engine.sync_engine (the synchronous path, matching the existing clean_in_fork pattern) and clear both async_engine and AsyncSession references.
When _configure_async_session() was extracted from configure_orm() in PR apache#51920, the corresponding cleanup in dispose_orm() was not updated. This left async_engine connections abandoned on process exit, gunicorn worker restarts, and atexit -- gradually exhausting PostgreSQL's max_connections. Dispose async_engine.sync_engine (the synchronous path, matching the existing clean_in_fork pattern) and clear both async_engine and AsyncSession references.
When _configure_async_session() was extracted from configure_orm() in PR #51920, the corresponding cleanup in dispose_orm() was not updated. This left async_engine connections abandoned on process exit, gunicorn worker restarts, and atexit -- gradually exhausting PostgreSQL's max_connections. Dispose async_engine.sync_engine (the synchronous path, matching the existing clean_in_fork pattern) and clear both async_engine and AsyncSession references.
When _configure_async_session() was extracted from configure_orm() in PR #51920, the corresponding cleanup in dispose_orm() was not updated. This left async_engine connections abandoned on process exit, gunicorn worker restarts, and atexit -- gradually exhausting PostgreSQL's max_connections. Dispose async_engine.sync_engine (the synchronous path, matching the existing clean_in_fork pattern) and clear both async_engine and AsyncSession references. (cherry picked from commit 69a88bf)
…5284) When _configure_async_session() was extracted from configure_orm() in PR #51920, the corresponding cleanup in dispose_orm() was not updated. This left async_engine connections abandoned on process exit, gunicorn worker restarts, and atexit -- gradually exhausting PostgreSQL's max_connections. Dispose async_engine.sync_engine (the synchronous path, matching the existing clean_in_fork pattern) and clear both async_engine and AsyncSession references. (cherry picked from commit 69a88bf) Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
…5284) When _configure_async_session() was extracted from configure_orm() in PR #51920, the corresponding cleanup in dispose_orm() was not updated. This left async_engine connections abandoned on process exit, gunicorn worker restarts, and atexit -- gradually exhausting PostgreSQL's max_connections. Dispose async_engine.sync_engine (the synchronous path, matching the existing clean_in_fork pattern) and clear both async_engine and AsyncSession references. (cherry picked from commit 69a88bf) Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
…5284) When _configure_async_session() was extracted from configure_orm() in PR #51920, the corresponding cleanup in dispose_orm() was not updated. This left async_engine connections abandoned on process exit, gunicorn worker restarts, and atexit -- gradually exhausting PostgreSQL's max_connections. Dispose async_engine.sync_engine (the synchronous path, matching the existing clean_in_fork pattern) and clear both async_engine and AsyncSession references. (cherry picked from commit 69a88bf) Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
…5284) When _configure_async_session() was extracted from configure_orm() in PR #51920, the corresponding cleanup in dispose_orm() was not updated. This left async_engine connections abandoned on process exit, gunicorn worker restarts, and atexit -- gradually exhausting PostgreSQL's max_connections. Dispose async_engine.sync_engine (the synchronous path, matching the existing clean_in_fork pattern) and clear both async_engine and AsyncSession references. (cherry picked from commit 69a88bf) Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
…5284) When _configure_async_session() was extracted from configure_orm() in PR #51920, the corresponding cleanup in dispose_orm() was not updated. This left async_engine connections abandoned on process exit, gunicorn worker restarts, and atexit -- gradually exhausting PostgreSQL's max_connections. Dispose async_engine.sync_engine (the synchronous path, matching the existing clean_in_fork pattern) and clear both async_engine and AsyncSession references. (cherry picked from commit 69a88bf) Co-authored-by: Kaxil Naik <kaxilnaik@gmail.com>
When _configure_async_session() was extracted from configure_orm() in PR apache#51920, the corresponding cleanup in dispose_orm() was not updated. This left async_engine connections abandoned on process exit, gunicorn worker restarts, and atexit -- gradually exhausting PostgreSQL's max_connections. Dispose async_engine.sync_engine (the synchronous path, matching the existing clean_in_fork pattern) and clear both async_engine and AsyncSession references.
Close#51711.
I initially wanted to just enhance the trigger endpoint to optionally stream until the run finishes, but it seems that FastAPI does not like this optionally stream idea. You can do it of course, but would loose a lot of the auto annotation reflection feature. So I opted to have a separate streaming endpoint instead.
This endpoint repeatedly emits a JSON object at the specified interval, until the dag reaches a finished state.
Tests to come.