diff --git a/airflow/www/static/js/trigger.js b/airflow/www/static/js/trigger.js index d2a0a4f687fd1..a235d5ee488a1 100644 --- a/airflow/www/static/js/trigger.js +++ b/airflow/www/static/js/trigger.js @@ -20,13 +20,32 @@ /* global document, CodeMirror, window */ const textArea = document.getElementById('json'); +const recentConfigList = document.getElementById('recent_configs'); const minHeight = 300; const maxHeight = window.innerHeight - 450; const height = maxHeight > minHeight ? maxHeight : minHeight; CodeMirror.fromTextArea(textArea, { lineNumbers: true, - mode: { name: 'javascript', json: true }, + mode: { + name: 'javascript', + json: true, + }, gutters: ['CodeMirror-lint-markers'], lint: true, -}).setSize(null, height); +}) + .setSize(null, height); + +function setRecentConfig(e) { + let { value } = e.target; + try { + const json = JSON.parse(value); + value = JSON.stringify(json, null, 2); + } catch (err) { + // eslint-disable-next-line no-console + console.error('config is not valid JSON format'); + } + document.querySelector('.CodeMirror').CodeMirror.setValue(value); +} + +recentConfigList.addEventListener('change', setRecentConfig); diff --git a/airflow/www/templates/airflow/trigger.html b/airflow/www/templates/airflow/trigger.html index 4a42990097e76..8c7e590aa2c55 100644 --- a/airflow/www/templates/airflow/trigger.html +++ b/airflow/www/templates/airflow/trigger.html @@ -44,12 +44,25 @@

Trigger DAG: {{ dag_id }}

- + +
+
+
+
+ +
- + +
+

To access configuration in your DAG use {{ '{{ dag_run.conf }}' }}. {% if is_dag_run_conf_overrides_params %} diff --git a/airflow/www/utils.py b/airflow/www/utils.py index 971a2c7a82261..4302c8d9f7e52 100644 --- a/airflow/www/utils.py +++ b/airflow/www/utils.py @@ -133,17 +133,24 @@ def get_mapped_summary(parent_instance, task_instances): } +def get_dag_run_conf(dag_run_conf: Any) -> tuple[str | None, bool]: + conf: str | None = None + + conf_is_json: bool = False + if isinstance(dag_run_conf, str): + conf = dag_run_conf + elif isinstance(dag_run_conf, (dict, list)) and any(dag_run_conf): + conf = json.dumps(dag_run_conf, sort_keys=True) + conf_is_json = True + + return conf, conf_is_json + + def encode_dag_run(dag_run: DagRun | None) -> dict[str, Any] | None: if not dag_run: return None - conf: str | None = None - conf_is_json: bool = False - if isinstance(dag_run.conf, str): - conf = dag_run.conf - elif isinstance(dag_run.conf, (dict, list)) and any(dag_run.conf): - conf = json.dumps(dag_run.conf, sort_keys=True) - conf_is_json = True + conf, conf_is_json = get_dag_run_conf(dag_run.conf) return { "run_id": dag_run.run_id, diff --git a/airflow/www/views.py b/airflow/www/views.py index a4126535cad64..81c54d7065946 100644 --- a/airflow/www/views.py +++ b/airflow/www/views.py @@ -1886,7 +1886,7 @@ def trigger(self, session=None): request_execution_date = request.values.get("execution_date", default=timezone.utcnow().isoformat()) is_dag_run_conf_overrides_params = conf.getboolean("core", "dag_run_conf_overrides_params") dag = get_airflow_app().dag_bag.get_dag(dag_id) - dag_orm = session.query(models.DagModel).filter(models.DagModel.dag_id == dag_id).first() + dag_orm = session.query(DagModel).filter(DagModel.dag_id == dag_id).first() if not dag_orm: flash(f"Cannot find dag {dag_id}") return redirect(origin) @@ -1895,6 +1895,22 @@ def trigger(self, session=None): flash(f"Cannot create dagruns because the dag {dag_id} has import errors", "error") return redirect(origin) + recent_runs = ( + session.query(DagRun.conf, func.max(DagRun.execution_date)) + .filter( + DagRun.dag_id == dag_id, + DagRun.run_type == DagRunType.MANUAL, + DagRun.conf.isnot(None), + ) + .group_by(DagRun.conf) + .order_by(func.max(DagRun.execution_date).desc()) + .limit(5) + .all() + ) + + recent_confs = [getattr(run, "conf") for run in recent_runs] + recent_confs = [conf[0] for conf in map(wwwutils.get_dag_run_conf, recent_confs) if conf[1]] + if request.method == "GET": # Populate conf textarea with conf requests parameter, or dag.params default_conf = "" @@ -1919,6 +1935,7 @@ def trigger(self, session=None): doc_md=doc_md, form=form, is_dag_run_conf_overrides_params=is_dag_run_conf_overrides_params, + recent_confs=recent_confs, ) try: @@ -1933,6 +1950,7 @@ def trigger(self, session=None): conf=request_conf, form=form, is_dag_run_conf_overrides_params=is_dag_run_conf_overrides_params, + recent_confs=recent_confs, ) dr = DagRun.find_duplicate(dag_id=dag_id, run_id=run_id, execution_date=execution_date) @@ -1966,6 +1984,7 @@ def trigger(self, session=None): conf=request_conf, form=form, is_dag_run_conf_overrides_params=is_dag_run_conf_overrides_params, + recent_confs=recent_confs, ) except json.decoder.JSONDecodeError: flash("Invalid JSON configuration, not parseable", "error") @@ -1977,6 +1996,7 @@ def trigger(self, session=None): conf=request_conf, form=form, is_dag_run_conf_overrides_params=is_dag_run_conf_overrides_params, + recent_confs=recent_confs, ) if unpause and dag.is_paused: