This is an automated email from the ASF dual-hosted git repository.
jedcunningham 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 85e648f970 Refactor DAG pages to be consistent (#25402)
85e648f970 is described below
commit 85e648f970f4ec63c7d3f5e96483204860b9fb26
Author: Jed Cunningham <[email protected]>
AuthorDate: Fri Jul 29 13:03:40 2022 -0700
Refactor DAG pages to be consistent (#25402)
There were a few DAG pages that were inconsistent:
- Missing error handling for non-existing dag_ids
- Not passing the orm model for the templates use (most notably missing
the "next run" info in the header)
---
airflow/www/views.py | 112 +++++++++++++++++++++++++++++----------------------
1 file changed, 64 insertions(+), 48 deletions(-)
diff --git a/airflow/www/views.py b/airflow/www/views.py
index da34997761..08ee4c6fa9 100644
--- a/airflow/www/views.py
+++ b/airflow/www/views.py
@@ -1237,31 +1237,27 @@ class Airflow(AirflowBaseView):
@provide_session
def code(self, dag_id, session=None):
"""Dag Code."""
- all_errors = ""
- dag_orm = None
+ dag = get_airflow_app().dag_bag.get_dag(dag_id, session=session)
+ dag_model = DagModel.get_dagmodel(dag_id, session=session)
+ if not dag:
+ flash(f'DAG "{dag_id}" seems to be missing.', "error")
+ return redirect(url_for('Airflow.index'))
+
+ wwwutils.check_import_errors(dag_model.fileloc, session)
+ wwwutils.check_dag_warnings(dag_model.dag_id, session)
try:
- dag_orm = DagModel.get_dagmodel(dag_id, session=session)
- code = DagCode.get_code_by_fileloc(dag_orm.fileloc)
+ code = DagCode.get_code_by_fileloc(dag_model.fileloc)
html_code = Markup(highlight(code, lexers.PythonLexer(),
HtmlFormatter(linenos=True)))
-
except Exception as e:
- all_errors += (
- "Exception encountered during "
- f"dag_id retrieval/dag retrieval fallback/code
highlighting:\n\n{e}\n"
- )
- html_code = Markup('<p>Failed to load DAG file
Code.</p><p>Details: {}</p>').format(
- escape(all_errors)
- )
-
- wwwutils.check_import_errors(dag_orm.fileloc, session)
- wwwutils.check_dag_warnings(dag_orm.dag_id, session)
+ error = f"Exception encountered during dag code retrieval/code
highlighting:\n\n{e}\n"
+ html_code = Markup('<p>Failed to load DAG file
Code.</p><p>Details: {}</p>').format(escape(error))
return self.render_template(
'airflow/dag_code.html',
html_code=html_code,
- dag=dag_orm,
- dag_model=dag_orm,
+ dag=dag,
+ dag_model=dag_model,
title=dag_id,
root=request.args.get('root'),
wrapped=conf.getboolean('webserver', 'default_wrap'),
@@ -1288,15 +1284,18 @@ class Airflow(AirflowBaseView):
@provide_session
def dag_details(self, dag_id, session=None):
"""Get Dag details."""
- dag = get_airflow_app().dag_bag.get_dag(dag_id)
+ dag = get_airflow_app().dag_bag.get_dag(dag_id, session=session)
dag_model = DagModel.get_dagmodel(dag_id, session=session)
-
- title = "DAG Details"
- root = request.args.get('root', '')
+ if not dag:
+ flash(f'DAG "{dag_id}" seems to be missing.', "error")
+ return redirect(url_for('Airflow.index'))
wwwutils.check_import_errors(dag.fileloc, session)
wwwutils.check_dag_warnings(dag.dag_id, session)
+ title = "DAG Details"
+ root = request.args.get('root', '')
+
states = (
session.query(TaskInstance.state,
sqla.func.count(TaskInstance.dag_id))
.filter(TaskInstance.dag_id == dag_id)
@@ -1336,6 +1335,7 @@ class Airflow(AirflowBaseView):
return self.render_template(
'airflow/dag_details.html',
dag=dag,
+ dag_model=dag_model,
title=title,
root=root,
states=states,
@@ -2692,8 +2692,8 @@ class Airflow(AirflowBaseView):
@provide_session
def grid(self, dag_id, session=None):
"""Get Dag's grid view."""
- dag = get_airflow_app().dag_bag.get_dag(dag_id)
- dag_model = DagModel.get_dagmodel(dag_id)
+ dag = get_airflow_app().dag_bag.get_dag(dag_id, session=session)
+ dag_model = DagModel.get_dagmodel(dag_id, session=session)
if not dag:
flash(f'DAG "{dag_id}" seems to be missing from DagBag.', "error")
return redirect(url_for('Airflow.index'))
@@ -2802,8 +2802,8 @@ class Airflow(AirflowBaseView):
else:
return func.date(column)
- dag = get_airflow_app().dag_bag.get_dag(dag_id)
- dag_model = DagModel.get_dagmodel(dag_id)
+ dag = get_airflow_app().dag_bag.get_dag(dag_id, session=session)
+ dag_model = DagModel.get_dagmodel(dag_id, session=session)
if not dag:
flash(f'DAG "{dag_id}" seems to be missing from DagBag.', "error")
return redirect(url_for('Airflow.index'))
@@ -2927,11 +2927,12 @@ class Airflow(AirflowBaseView):
@provide_session
def graph(self, dag_id, session=None):
"""Get DAG as Graph."""
- dag = get_airflow_app().dag_bag.get_dag(dag_id)
- dag_model = DagModel.get_dagmodel(dag_id)
+ dag = get_airflow_app().dag_bag.get_dag(dag_id, session=session)
+ dag_model = DagModel.get_dagmodel(dag_id, session=session)
if not dag:
flash(f'DAG "{dag_id}" seems to be missing.', "error")
return redirect(url_for('Airflow.index'))
+
wwwutils.check_import_errors(dag.fileloc, session)
wwwutils.check_dag_warnings(dag.dag_id, session)
@@ -3037,17 +3038,16 @@ class Airflow(AirflowBaseView):
@provide_session
def duration(self, dag_id, session=None):
"""Get Dag as duration graph."""
- default_dag_run = conf.getint('webserver',
'default_dag_run_display_number')
- dag_model = DagModel.get_dagmodel(dag_id)
-
- dag: Optional[DAG] = get_airflow_app().dag_bag.get_dag(dag_id)
- if dag is None:
+ dag = get_airflow_app().dag_bag.get_dag(dag_id, session=session)
+ dag_model = DagModel.get_dagmodel(dag_id, session=session)
+ if not dag:
flash(f'DAG "{dag_id}" seems to be missing.', "error")
return redirect(url_for('Airflow.index'))
wwwutils.check_import_errors(dag.fileloc, session)
wwwutils.check_dag_warnings(dag.dag_id, session)
+ default_dag_run = conf.getint('webserver',
'default_dag_run_display_number')
base_date = request.args.get('base_date')
num_runs = request.args.get('num_runs', default=default_dag_run,
type=int)
@@ -3192,9 +3192,16 @@ class Airflow(AirflowBaseView):
@provide_session
def tries(self, dag_id, session=None):
"""Shows all tries."""
+ dag = get_airflow_app().dag_bag.get_dag(dag_id, session=session)
+ dag_model = DagModel.get_dagmodel(dag_id, session=session)
+ if not dag:
+ flash(f'DAG "{dag_id}" seems to be missing.', "error")
+ return redirect(url_for('Airflow.index'))
+
+ wwwutils.check_import_errors(dag.fileloc, session)
+ wwwutils.check_dag_warnings(dag.dag_id, session)
+
default_dag_run = conf.getint('webserver',
'default_dag_run_display_number')
- dag = get_airflow_app().dag_bag.get_dag(dag_id)
- dag_model = DagModel.get_dagmodel(dag_id)
base_date = request.args.get('base_date')
num_runs = request.args.get('num_runs', default=default_dag_run,
type=int)
@@ -3203,9 +3210,6 @@ class Airflow(AirflowBaseView):
else:
base_date = dag.get_latest_execution_date() or timezone.utcnow()
- wwwutils.check_import_errors(dag.fileloc, session)
- wwwutils.check_dag_warnings(dag.dag_id, session)
-
root = request.args.get('root')
if root:
dag = dag.partial_subset(task_ids_or_regex=root,
include_upstream=True, include_downstream=False)
@@ -3283,9 +3287,16 @@ class Airflow(AirflowBaseView):
@provide_session
def landing_times(self, dag_id, session=None):
"""Shows landing times."""
+ dag = get_airflow_app().dag_bag.get_dag(dag_id, session=session)
+ dag_model = DagModel.get_dagmodel(dag_id, session=session)
+ if not dag:
+ flash(f'DAG "{dag_id}" seems to be missing.', "error")
+ return redirect(url_for('Airflow.index'))
+
+ wwwutils.check_import_errors(dag.fileloc, session)
+ wwwutils.check_dag_warnings(dag.dag_id, session)
+
default_dag_run = conf.getint('webserver',
'default_dag_run_display_number')
- dag: DAG = get_airflow_app().dag_bag.get_dag(dag_id)
- dag_model = DagModel.get_dagmodel(dag_id)
base_date = request.args.get('base_date')
num_runs = request.args.get('num_runs', default=default_dag_run,
type=int)
@@ -3294,9 +3305,6 @@ class Airflow(AirflowBaseView):
else:
base_date = dag.get_latest_execution_date() or timezone.utcnow()
- wwwutils.check_import_errors(dag.fileloc, session)
- wwwutils.check_dag_warnings(dag.dag_id, session)
-
root = request.args.get('root')
if root:
dag = dag.partial_subset(task_ids_or_regex=root,
include_upstream=True, include_downstream=False)
@@ -3403,16 +3411,19 @@ class Airflow(AirflowBaseView):
@provide_session
def gantt(self, dag_id, session=None):
"""Show GANTT chart."""
- dag = get_airflow_app().dag_bag.get_dag(dag_id)
- dag_model = DagModel.get_dagmodel(dag_id)
+ dag = get_airflow_app().dag_bag.get_dag(dag_id, session=session)
+ dag_model = DagModel.get_dagmodel(dag_id, session=session)
+ if not dag:
+ flash(f'DAG "{dag_id}" seems to be missing.', "error")
+ return redirect(url_for('Airflow.index'))
+
+ wwwutils.check_import_errors(dag.fileloc, session)
+ wwwutils.check_dag_warnings(dag.dag_id, session)
root = request.args.get('root')
if root:
dag = dag.partial_subset(task_ids_or_regex=root,
include_upstream=True, include_downstream=False)
- wwwutils.check_import_errors(dag.fileloc, session)
- wwwutils.check_dag_warnings(dag.dag_id, session)
-
dt_nr_dr_data = get_date_time_num_runs_dag_runs_form_data(request,
session, dag)
dttm = dt_nr_dr_data['dttm']
dag_run = dag.get_dagrun(execution_date=dttm)
@@ -3680,7 +3691,11 @@ class Airflow(AirflowBaseView):
@provide_session
def audit_log(self, session=None):
dag_id = request.args.get('dag_id')
- dag = get_airflow_app().dag_bag.get_dag(dag_id)
+ dag = get_airflow_app().dag_bag.get_dag(dag_id, session=session)
+ dag_model = DagModel.get_dagmodel(dag_id, session=session)
+ if not dag:
+ flash(f'DAG "{dag_id}" seems to be missing from DagBag.', "error")
+ return redirect(url_for('Airflow.index'))
included_events = conf.get('webserver', 'audit_view_included_events',
fallback=None)
excluded_events = conf.get('webserver', 'audit_view_excluded_events',
fallback=None)
@@ -3697,6 +3712,7 @@ class Airflow(AirflowBaseView):
content = self.render_template(
'airflow/dag_audit_log.html',
dag=dag,
+ dag_model=dag_model,
root=request.args.get('root'),
dag_id=dag_id,
dag_logs=dag_audit_logs,