1#
2# Licensed to the Apache Software Foundation (ASF) under one
3# or more contributor license agreements. See the NOTICE file
4# distributed with this work for additional information
5# regarding copyright ownership. The ASF licenses this file
6# to you under the Apache License, Version 2.0 (the
7# "License"); you may not use this file except in compliance
8# with the License. You may obtain a copy of the License at
9#
10# http://www.apache.org/licenses/LICENSE-2.0
11#
12# Unless required by applicable law or agreed to in writing,
13# software distributed under the License is distributed on an
14# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
15# KIND, either express or implied. See the License for the
16# specific language governing permissions and limitations
17# under the License.
18"""Base configuration parser with pure parsing logic."""
19
20from __future__ import annotations
21
22import contextlib
23import datetime
24import functools
25import itertools
26import json
27import logging
28import os
29import shlex
30import subprocess
31import sys
32import warnings
33from collections.abc import Callable, Generator, Iterable
34from configparser import ConfigParser, NoOptionError, NoSectionError
35from contextlib import contextmanager
36from copy import deepcopy
37from enum import Enum
38from json.decoder import JSONDecodeError
39from re import Pattern
40from typing import IO, TYPE_CHECKING, Any, TypeVar, overload
41
42from .exceptions import AirflowConfigException
43
44log = logging.getLogger(__name__)
45
46
47def _build_kwarg_env_prefix(section: str, kwargs_key: str) -> str:
48 """
49 Build env prefix for per-key backend kwargs.
50
51 ("secrets", "backend_kwargs") -> "AIRFLOW__SECRETS__BACKEND_KWARG__"
52 ("workers", "secrets_backend_kwargs") -> "AIRFLOW__WORKERS__SECRETS_BACKEND_KWARG__"
53 """
54 singular_key = kwargs_key.replace("_kwargs", "_kwarg")
55 return f"{ENV_VAR_PREFIX}{section.upper()}__{singular_key.upper()}__"
56
57
58def _collect_kwarg_env_vars(prefix: str) -> dict[str, str]:
59 """
60 Scan os.environ for per-key secrets backend kwargs.
61
62 AIRFLOW__SECRETS__BACKEND_KWARG__ROLE_ID -> {"role_id": value}
63 Values are raw strings (not JSON-parsed).
64 Empty keys (trailing __ with no suffix) are ignored.
65 """
66 overrides: dict[str, str] = {}
67 for env_var, value in os.environ.items():
68 if env_var.startswith(prefix):
69 kwarg_key = env_var[len(prefix) :].lower()
70 if kwarg_key:
71 overrides[kwarg_key] = value
72 return overrides
73
74
75ConfigType = str | int | float | bool
76ConfigOptionsDictType = dict[str, ConfigType]
77ConfigSectionSourcesType = dict[str, str | tuple[str, str]]
78ConfigSourcesType = dict[str, ConfigSectionSourcesType]
79ENV_VAR_PREFIX = "AIRFLOW__"
80# Separates the team name from the base section name in a team scoped config file section.
81TEAM_SECTION_SEPARATOR = "="
82# (section, key) pairs that may hold a custom secrets backend class. _get_custom_secret_backend()
83# reads this to pick the pair for the current mode; callers that only need to know whether *some*
84# backend is configured (e.g. a CLI warning), without instantiating one, read all of the values.
85# Named without "secret" so static analysis (e.g. CodeQL's clear-text-logging query) doesn't treat
86# the (non-secret) section/key strings sourced from here as sensitive data.
87CUSTOM_BACKEND_CONFIG_KEYS = {
88 "general": ("secrets", "backend"),
89 "worker": ("workers", "secrets_backend"),
90}
91
92
93def team_section_name(team_name: str, section: str) -> str:
94 """
95 Build the config file section name that holds the team scoped overrides of ``section``.
96
97 :param team_name: name of the team the overrides belong to
98 :param section: base section name that is being overridden
99 :return: the team scoped section name, e.g. ``team_a=celery``
100 """
101 return f"{team_name}{TEAM_SECTION_SEPARATOR}{section}"
102
103
104def base_section_name(section: str) -> str:
105 """
106 Return the base section name of a possibly team scoped config file section.
107
108 Team scoped sections are built by :func:`team_section_name`. Base section names never contain
109 the separator, so the name is split on the last one - that way the base section is recovered
110 even for a team name that contains the separator itself.
111
112 :param section: section name, either a base one or a team scoped one
113 :return: the base section name, which is ``section`` itself when it is not team scoped
114 """
115 _, separator, base_section = section.rpartition(TEAM_SECTION_SEPARATOR)
116 return base_section if separator else section
117
118
119if TYPE_CHECKING:
120 from airflow.providers_manager import ProvidersManager
121 from airflow.sdk.providers_manager_runtime import ProvidersManagerTaskRuntime
122
123
124class ValueNotFound:
125 """Object of this is raised when a configuration value cannot be found."""
126
127 pass
128
129
130VALUE_NOT_FOUND_SENTINEL = ValueNotFound()
131
132
133@overload
134def expand_env_var(env_var: None) -> None: ...
135@overload
136def expand_env_var(env_var: str) -> str: ...
137
138
139def expand_env_var(env_var: str | None) -> str | None:
140 """
141 Expand (potentially nested) env vars.
142
143 Repeat and apply `expandvars` and `expanduser` until
144 interpolation stops having any effect.
145 """
146 if not env_var or not isinstance(env_var, str):
147 return env_var
148 while True:
149 interpolated = os.path.expanduser(os.path.expandvars(str(env_var)))
150 if interpolated == env_var:
151 return interpolated
152 env_var = interpolated
153
154
155def run_command(command: str) -> str:
156 """Run command and returns stdout."""
157 process = subprocess.Popen(
158 shlex.split(command), stdout=subprocess.PIPE, stderr=subprocess.PIPE, close_fds=True
159 )
160 output, stderr = (stream.decode(sys.getdefaultencoding(), "ignore") for stream in process.communicate())
161
162 if process.returncode != 0:
163 raise AirflowConfigException(
164 f"Cannot execute {command}. Error code is: {process.returncode}. "
165 f"Output: {output}, Stderr: {stderr}"
166 )
167
168 return output
169
170
171def _is_template(configuration_description: dict[str, dict[str, Any]], section: str, key: str) -> bool:
172 """
173 Check if the config is a template.
174
175 :param configuration_description: description of configuration
176 :param section: section
177 :param key: key
178 :return: True if the config is a template
179 """
180 return configuration_description.get(section, {}).get(key, {}).get("is_template", False)
181
182
183def configure_parser_from_configuration_description(
184 parser: ConfigParser,
185 configuration_description: dict[str, dict[str, Any]],
186 all_vars: dict[str, Any],
187) -> None:
188 """
189 Configure a ConfigParser based on configuration description.
190
191 :param parser: ConfigParser to configure
192 :param configuration_description: configuration description from config.yml
193 """
194 for section, section_desc in configuration_description.items():
195 parser.add_section(section)
196 options = section_desc["options"]
197 for key in options:
198 default_value = options[key]["default"]
199 is_template = options[key].get("is_template", False)
200 if (default_value is not None) and not (
201 options[key].get("version_deprecated") or options[key].get("deprecation_reason")
202 ):
203 if is_template or not isinstance(default_value, str):
204 parser.set(section, key, str(default_value))
205 else:
206 try:
207 parser.set(section, key, default_value.format(**all_vars))
208 except (KeyError, ValueError):
209 parser.set(section, key, default_value)
210
211
212def create_provider_cfg_config_fallback_defaults(
213 provider_config_fallback_defaults_cfg_path: str,
214) -> ConfigParser:
215 """
216 Create fallback defaults for configuration.
217
218 This parser contains provider defaults for Airflow configuration, containing fallback default values
219 that might be needed when provider classes are being imported - before provider's configuration
220 is loaded.
221
222 Unfortunately airflow currently performs a lot of stuff during importing and some of that might lead
223 to retrieving provider configuration before the defaults for the provider are loaded.
224
225 Those are only defaults, so if you have "real" values configured in your configuration (.cfg file or
226 environment variables) those will be used as usual.
227
228 NOTE!! Do NOT attempt to remove those default fallbacks thinking that they are unnecessary duplication,
229 at least not until we fix the way how airflow imports "do stuff". This is unlikely to succeed.
230
231 You've been warned!
232
233 :param provider_config_fallback_defaults_cfg_path: path to the provider config fallback defaults .cfg file
234 """
235 config_parser = ConfigParser()
236 config_parser.read(provider_config_fallback_defaults_cfg_path)
237 return config_parser
238
239
240class AirflowConfigParser(ConfigParser):
241 """
242 Base configuration parser with pure parsing logic.
243
244 This class provides the core parsing methods that work with:
245 - configuration_description: dict describing config options (required in __init__)
246 - _default_values: ConfigParser with default values (required in __init__)
247 - deprecated_options: class attribute mapping new -> old options
248 - deprecated_sections: class attribute mapping new -> old sections
249 """
250
251 # A mapping of section -> setting -> { old, replace } for deprecated default values.
252 # Subclasses can override this to define deprecated values that should be upgraded.
253 deprecated_values: dict[str, dict[str, tuple[Pattern, str]]] = {}
254
255 # A mapping of (new section, new option) -> (old section, old option, since_version).
256 # When reading new option, the old option will be checked to see if it exists. If it does a
257 # DeprecationWarning will be issued and the old option will be used instead
258 deprecated_options: dict[tuple[str, str], tuple[str, str, str]] = {
259 ("dag_processor", "dag_file_processor_timeout"): ("core", "dag_file_processor_timeout", "3.0"),
260 ("dag_processor", "refresh_interval"): ("scheduler", "dag_dir_list_interval", "3.0"),
261 ("api", "base_url"): ("webserver", "base_url", "3.0"),
262 ("api", "host"): ("webserver", "web_server_host", "3.0"),
263 ("api", "port"): ("webserver", "web_server_port", "3.0"),
264 ("api", "workers"): ("webserver", "workers", "3.0"),
265 ("api", "worker_timeout"): ("webserver", "web_server_worker_timeout", "3.0"),
266 ("api", "ssl_cert"): ("webserver", "web_server_ssl_cert", "3.0"),
267 ("api", "ssl_key"): ("webserver", "web_server_ssl_key", "3.0"),
268 ("api", "access_logfile"): ("webserver", "access_logfile", "3.0"),
269 ("triggerer", "capacity"): ("triggerer", "default_capacity", "3.0"),
270 ("api", "expose_config"): ("webserver", "expose_config", "3.0.1"),
271 ("fab", "access_denied_message"): ("webserver", "access_denied_message", "3.0.2"),
272 ("fab", "expose_hostname"): ("webserver", "expose_hostname", "3.0.2"),
273 ("fab", "navbar_color"): ("webserver", "navbar_color", "3.0.2"),
274 ("fab", "navbar_text_color"): ("webserver", "navbar_text_color", "3.0.2"),
275 ("fab", "navbar_hover_color"): ("webserver", "navbar_hover_color", "3.0.2"),
276 ("fab", "navbar_text_hover_color"): ("webserver", "navbar_text_hover_color", "3.0.2"),
277 ("api", "secret_key"): ("webserver", "secret_key", "3.0.2"),
278 ("api", "enable_swagger_ui"): ("webserver", "enable_swagger_ui", "3.0.2"),
279 ("dag_processor", "parsing_pre_import_modules"): ("scheduler", "parsing_pre_import_modules", "3.0.4"),
280 ("api", "grid_view_sorting_order"): ("webserver", "grid_view_sorting_order", "3.1.0"),
281 ("api", "log_fetch_timeout_sec"): ("webserver", "log_fetch_timeout_sec", "3.1.0"),
282 ("api", "hide_paused_dags_by_default"): ("webserver", "hide_paused_dags_by_default", "3.1.0"),
283 ("core", "num_dag_runs_to_retain_rendered_fields"): (
284 "core",
285 "max_num_rendered_ti_fields_per_task",
286 "3.2.0",
287 ),
288 ("api", "page_size"): ("webserver", "page_size", "3.1.0"),
289 ("api", "default_wrap"): ("webserver", "default_wrap", "3.1.0"),
290 ("api", "auto_refresh_interval"): ("webserver", "auto_refresh_interval", "3.1.0"),
291 ("api", "require_confirmation_dag_change"): ("webserver", "require_confirmation_dag_change", "3.1.0"),
292 ("api", "instance_name"): ("webserver", "instance_name", "3.1.0"),
293 ("api", "log_config"): ("api", "access_logfile", "3.1.0"),
294 ("scheduler", "ti_metrics_interval"): ("scheduler", "running_metrics_interval", "3.2.0"),
295 ("api", "fallback_page_limit"): ("api", "page_size", "3.2.0"),
296 ("workers", "missing_dag_retries"): ("workers", "missing_dag_retires", "3.1.8"),
297 ("core", "execution_api_server_url"): ("workers", "execution_api_server_url", "3.0"),
298 ("database", "sql_alchemy_conn"): ("core", "sql_alchemy_conn", "3.0"),
299 }
300
301 # A mapping of new section -> (old section, since_version).
302 deprecated_sections: dict[str, tuple[str, str]] = {}
303
304 @property
305 def _lookup_sequence(self) -> list[Callable]:
306 """
307 Define the sequence of lookup methods for get(). The definition here does not have provider lookup.
308
309 Subclasses can override this to customise lookup order.
310 """
311 lookup_methods = [
312 self._get_environment_variables,
313 self._get_option_from_config_file,
314 self._get_option_from_commands,
315 self._get_option_from_secrets,
316 self._get_option_from_defaults,
317 ]
318 if self._use_providers_configuration:
319 # Provider fallback lookups are last so they have the lowest priority in the lookup sequence.
320 lookup_methods += [
321 self._get_option_from_provider_metadata_config_fallbacks,
322 self._get_option_from_provider_cfg_config_fallbacks,
323 ]
324 return lookup_methods
325
326 @functools.cached_property
327 def configuration_description(self) -> dict[str, dict[str, Any]]:
328 """
329 Return configuration description from multiple sources.
330
331 Respects the ``_use_providers_configuration`` flag to decide whether to include
332 provider configuration.
333
334 The merged description is built as follows:
335
336 1. Start from the base configuration description provided in ``__init__``, usually
337 loaded from ``config.yml`` in core. Values defined here are never overridden.
338 2. Merge provider metadata from ``_provider_metadata_configuration_description``,
339 loaded from provider packages' ``get_provider_info`` method. Only adds missing
340 sections/options; does not overwrite existing entries from the base configuration.
341 3. Merge default values from ``_provider_cfg_config_fallback_default_values``,
342 loaded from ``provider_config_fallback_defaults.cfg``. Only sets ``"default"``
343 (and heuristically ``"sensitive"``) for options that do not already define them.
344
345 Base configuration takes precedence, then provider metadata fills in missing
346 descriptions/options, and finally cfg-based fallbacks provide defaults only where
347 none are defined.
348
349 We use ``cached_property`` to cache the merged result; clear this cache (via
350 ``invalidate_cache``) when toggling ``_use_providers_configuration``.
351 """
352 if not self._use_providers_configuration:
353 return self._configuration_description
354
355 merged_description: dict[str, dict[str, Any]] = deepcopy(self._configuration_description)
356
357 # Merge full provider config descriptions (with metadata like sensitive, description, etc.)
358 # from provider packages' get_provider_info method, reusing the cached raw dict.
359 for section, section_content in self._provider_metadata_configuration_description.items():
360 if section not in merged_description:
361 merged_description[section] = deepcopy(section_content)
362 else:
363 existing_options = merged_description[section].setdefault("options", {})
364 for option, option_content in section_content.get("options", {}).items():
365 if option not in existing_options:
366 existing_options[option] = deepcopy(option_content)
367
368 # Merge default values from cfg-based fallbacks (key=value only, no metadata).
369 # Uses setdefault so provider metadata values above take priority.
370 cfg = self._provider_cfg_config_fallback_default_values
371 for section in cfg.sections():
372 section_options = merged_description.setdefault(section, {"options": {}}).setdefault(
373 "options", {}
374 )
375 for option in cfg.options(section):
376 opt_dict = section_options.setdefault(option, {})
377 opt_dict.setdefault("default", cfg.get(section, option))
378 # For cfg-only options with no provider metadata, infer sensitivity from name.
379 if "sensitive" not in opt_dict and option.endswith(("password", "secret")):
380 opt_dict["sensitive"] = True
381
382 return merged_description
383
384 @property
385 def _config_sources_for_as_dict(self) -> list[tuple[str, ConfigParser]]:
386 """Override the base method to add provider fallbacks when providers are loaded."""
387 sources: list[tuple[str, ConfigParser]] = []
388 if self._use_providers_configuration:
389 # Provider fallback defaults are listed first so they have the lowest priority
390 # in as_dict()'s "last source wins" semantics.
391 sources += [
392 ("provider-cfg-fallback-defaults", self._provider_cfg_config_fallback_default_values),
393 (
394 "provider-metadata-fallback-defaults",
395 self._provider_metadata_config_fallback_default_values,
396 ),
397 ]
398 sources += [
399 ("default", self._default_values),
400 ("airflow.cfg", self),
401 ]
402 return sources
403
404 def _get_option_from_provider_cfg_config_fallbacks(
405 self,
406 deprecated_key: str | None,
407 deprecated_section: str | None,
408 key: str,
409 section: str,
410 issue_warning: bool = True,
411 extra_stacklevel: int = 0,
412 **kwargs,
413 ) -> str | ValueNotFound:
414 """Get config option from provider fallback defaults."""
415 value = self.get_from_provider_cfg_config_fallback_defaults(section, key, **kwargs)
416 if value is not VALUE_NOT_FOUND_SENTINEL:
417 return value
418 return VALUE_NOT_FOUND_SENTINEL
419
420 def _get_option_from_provider_metadata_config_fallbacks(
421 self,
422 deprecated_key: str | None,
423 deprecated_section: str | None,
424 key: str,
425 section: str,
426 issue_warning: bool = True,
427 extra_stacklevel: int = 0,
428 **kwargs,
429 ) -> str | ValueNotFound:
430 """Get config option from provider metadata fallback defaults."""
431 value = self.get_from_provider_metadata_config_fallback_defaults(section, key, **kwargs)
432 if value is not VALUE_NOT_FOUND_SENTINEL:
433 return value
434 return VALUE_NOT_FOUND_SENTINEL
435
436 def get_from_provider_cfg_config_fallback_defaults(self, section: str, key: str, **kwargs) -> Any:
437 """Get provider config fallback default values."""
438 raw = kwargs.get("raw", False)
439 vars_ = kwargs.get("vars")
440 return self._provider_cfg_config_fallback_default_values.get(
441 section, key, fallback=VALUE_NOT_FOUND_SENTINEL, raw=raw, vars=vars_
442 )
443
444 @functools.cached_property
445 def _provider_metadata_configuration_description(self) -> dict[str, dict[str, Any]]:
446 """Raw provider configuration descriptions with full metadata (sensitive, description, etc.)."""
447 result: dict[str, dict[str, Any]] = {}
448 for _, config in self._provider_manager_type().provider_configs:
449 result.update(config)
450 return result
451
452 @functools.cached_property
453 def _provider_metadata_config_fallback_default_values(self) -> ConfigParser:
454 """Return Provider metadata config fallback default values."""
455 return self._create_default_config_parser_callable(self._provider_metadata_configuration_description)
456
457 def get_from_provider_metadata_config_fallback_defaults(self, section: str, key: str, **kwargs) -> Any:
458 """Get provider metadata config fallback default values."""
459 raw = kwargs.get("raw", False)
460 vars_ = kwargs.get("vars")
461 return self._provider_metadata_config_fallback_default_values.get(
462 section, key, fallback=VALUE_NOT_FOUND_SENTINEL, raw=raw, vars=vars_
463 )
464
465 @property
466 def _validators(self) -> list[Callable[[], None]]:
467 """
468 Return list of validators defined on a config parser class. Base class will return an empty list.
469
470 Subclasses can override this to customize the validators that are run during validation on the
471 config parser instance.
472 """
473 return []
474
475 def validate(self) -> None:
476 """Run all registered validators."""
477 for validator in self._validators:
478 validator()
479 self.is_validated = True
480
481 def _validate_deprecated_values(self) -> None:
482 """Validate and upgrade deprecated default values."""
483 for section, replacement in self.deprecated_values.items():
484 for name, info in replacement.items():
485 old, new = info
486 current_value = self.get(section, name, fallback="")
487 if self._using_old_value(old, current_value):
488 self.upgraded_values[(section, name)] = current_value
489 new_value = old.sub(new, current_value)
490 self._update_env_var(section=section, name=name, new_value=new_value)
491 self._create_future_warning(
492 name=name,
493 section=section,
494 current_value=current_value,
495 new_value=new_value,
496 )
497
498 def _using_old_value(self, old: Pattern, current_value: str) -> bool:
499 """Check if current_value matches the old pattern."""
500 return old.search(current_value) is not None
501
502 def _update_env_var(self, section: str, name: str, new_value: str) -> None:
503 """Update environment variable with new value."""
504 env_var = self._env_var_name(section, name)
505 # Set it as an env var so that any subprocesses keep the same override!
506 os.environ[env_var] = new_value
507
508 @staticmethod
509 def _create_future_warning(name: str, section: str, current_value: Any, new_value: Any) -> None:
510 """Create a FutureWarning for deprecated default values."""
511 warnings.warn(
512 f"The {name!r} setting in [{section}] has the old default value of {current_value!r}. "
513 f"This value has been changed to {new_value!r} in the running config, but please update your config.",
514 FutureWarning,
515 stacklevel=3,
516 )
517
518 def __init__(
519 self,
520 configuration_description: dict[str, dict[str, Any]],
521 _default_values: ConfigParser,
522 provider_manager_type: type[ProvidersManager] | type[ProvidersManagerTaskRuntime],
523 create_default_config_parser_callable: Callable[[dict[str, dict[str, Any]]], ConfigParser],
524 provider_config_fallback_defaults_cfg_path: str,
525 *args,
526 **kwargs,
527 ):
528 """
529 Initialize the parser.
530
531 :param configuration_description: Description of configuration options
532 :param _default_values: ConfigParser with default values
533 :param provider_manager_type: Either ProvidersManager or ProvidersManagerTaskRuntime, depending on the context of the caller.
534 :param create_default_config_parser_callable: The `create_default_config_parser` function from core or SDK, depending on the context of the caller.
535 :param provider_config_fallback_defaults_cfg_path: Path to the `provider_config_fallback_defaults.cfg` file.
536 """
537 super().__init__(*args, **kwargs)
538 self._configuration_description = configuration_description
539 self._default_values = _default_values
540 self._provider_manager_type = provider_manager_type
541 self._create_default_config_parser_callable = create_default_config_parser_callable
542 self._provider_cfg_config_fallback_default_values = create_provider_cfg_config_fallback_defaults(
543 provider_config_fallback_defaults_cfg_path
544 )
545 self._suppress_future_warnings = False
546 self.upgraded_values: dict[tuple[str, str], str] = {}
547 # The _use_providers_configuration flag will always be True unless we call `write(include_providers=False)` or `with self.make_sure_configuration_loaded(with_providers=False)`.
548 # Even when we call those methods, the flag will be set back to True after the method is done, so it only affects the current call to `as_dict()` and does not have any effect on subsequent calls.
549 self._use_providers_configuration = True
550
551 def invalidate_cache(self) -> None:
552 """
553 Clear all ``functools.cached_property`` entries on this instance.
554
555 Call this after mutating class-level attributes (e.g. ``deprecated_options``)
556 so that derived cached properties are recomputed on next access.
557 """
558 for attr_name in (
559 name
560 for name in dir(type(self))
561 if isinstance(getattr(type(self), name, None), functools.cached_property)
562 ):
563 self.__dict__.pop(attr_name, None)
564
565 def _invalidate_provider_flag_caches(self) -> None:
566 """Invalidate caches related to provider configuration flags."""
567 self.__dict__.pop("configuration_description", None)
568 self.__dict__.pop("sensitive_config_values", None)
569
570 @functools.cached_property
571 def inversed_deprecated_options(self):
572 """Build inverse mapping from old options to new options."""
573 return {(sec, name): key for key, (sec, name, ver) in self.deprecated_options.items()}
574
575 @functools.cached_property
576 def inversed_deprecated_sections(self):
577 """Build inverse mapping from old sections to new sections."""
578 return {
579 old_section: new_section for new_section, (old_section, ver) in self.deprecated_sections.items()
580 }
581
582 @functools.cached_property
583 def sensitive_config_values(self) -> set[tuple[str, str]]:
584 """Get set of sensitive config values that should be masked."""
585 flattened = {
586 (s, k): item
587 for s, s_c in self.configuration_description.items()
588 for k, item in s_c.get("options", {}).items()
589 }
590 sensitive = {
591 (section.lower(), key.lower())
592 for (section, key), v in flattened.items()
593 if v.get("sensitive") is True
594 }
595 depr_option = {self.deprecated_options[x][:-1] for x in sensitive if x in self.deprecated_options}
596 depr_section = {
597 (self.deprecated_sections[s][0], k) for s, k in sensitive if s in self.deprecated_sections
598 }
599 sensitive.update(depr_section, depr_option)
600 return sensitive
601
602 def _names_sensitive_team_env_var(self, env_var: str) -> bool:
603 """
604 Check whether an environment variable name is a team scoped override of a sensitive option.
605
606 Team scoped variables are named ``AIRFLOW__<TEAM>___<SECTION>__<KEY>`` (see
607 :meth:`_env_var_name`). A team name may contain underscores itself, so the team name is not
608 parsed out of the variable name; the name is matched against the tail that each option
609 registered as sensitive contributes instead. The ``_CMD`` / ``_SECRET`` fallbacks are
610 matched too, because - unlike for a base section - they are never resolved into their value.
611
612 :param env_var: environment variable name
613 :return: True if the variable holds a team scoped value of an option registered as sensitive
614 """
615 env_var = env_var.upper()
616 if not env_var.startswith(ENV_VAR_PREFIX):
617 return False
618 # Every tail matched below starts with the ``___`` separating the team name from the
619 # section, so a name not containing it cannot be a team scoped one.
620 if "___" not in env_var:
621 return False
622 for section, key in self.sensitive_config_values:
623 option_tail = self._env_var_name(section, key).removeprefix(ENV_VAR_PREFIX)
624 for tail in (f"___{option_tail}", f"___{option_tail}_CMD", f"___{option_tail}_SECRET"):
625 # The team name sits between the prefix and the tail, so it must not be empty.
626 if env_var.endswith(tail) and len(env_var) > len(ENV_VAR_PREFIX) + len(tail):
627 return True
628 return False
629
630 def is_sensitive_option(self, section: str, key: str) -> bool:
631 """
632 Check whether the value of ``key`` in ``section`` is registered as sensitive.
633
634 Options are registered as sensitive under their base section name, while a team scoped
635 override of the very same option is held by a ``<team name>=<section>`` config file section
636 or by an ``AIRFLOW__<TEAM>___<SECTION>__<KEY>`` environment variable. Both of those
637 spellings are resolved back to the base option here, so that a team scoped value is treated
638 exactly like the base one. A name that does not resolve to a registered option is not
639 sensitive - so this only ever recognises more options as sensitive, never fewer.
640
641 :param section: section name, either a base one or a team scoped one
642 :param key: option name
643 :return: True if the value of the option should be treated as sensitive
644 """
645 section = section.lower()
646 key = key.lower()
647 if (section, key) in self.sensitive_config_values:
648 return True
649 base_section = base_section_name(section)
650 if base_section != section:
651 if (base_section, key) in self.sensitive_config_values:
652 return True
653 # A team scoped ``_cmd`` / ``_secret`` fallback is not resolved into its value, so it
654 # stays in the output as configured and has to be recognised on its own.
655 for fallback_suffix in ("_cmd", "_secret"):
656 if not key.endswith(fallback_suffix):
657 continue
658 if (base_section, key.removesuffix(fallback_suffix)) in self.sensitive_config_values:
659 return True
660 # A team scoped environment variable is reported under the section and key its name splits
661 # into, which is neither the base nor the team scoped section name.
662 return self._names_sensitive_team_env_var(self._env_var_name(section, key))
663
664 def _update_defaults_from_string(self, config_string: str) -> None:
665 """
666 Update the defaults in _default_values based on values in config_string ("ini" format).
667
668 Override shared parser's method to add validation for template variables.
669 Note that those values are not validated and cannot contain variables because we are using
670 regular config parser to load them. This method is used to test the config parser in unit tests.
671
672 :param config_string: ini-formatted config string
673 """
674 parser = ConfigParser()
675 parser.read_string(config_string)
676 for section in parser.sections():
677 if section not in self._default_values.sections():
678 self._default_values.add_section(section)
679 errors = False
680 for key, value in parser.items(section):
681 if not self.is_template(section, key) and "{" in value:
682 errors = True
683 log.error(
684 "The %s.%s value %s read from string contains variable. This is not supported",
685 section,
686 key,
687 value,
688 )
689 self._default_values.set(section, key, value)
690 if errors:
691 raise AirflowConfigException(
692 f"The string config passed as default contains variables. "
693 f"This is not supported. String config: {config_string}"
694 )
695
696 def get_default_value(self, section: str, key: str, fallback: Any = None, raw=False, **kwargs) -> Any:
697 """
698 Retrieve default value from default config parser, including provider fallbacks.
699
700 This will retrieve the default value from the core default config parser first. If not found
701 and providers configuration is loaded, it also checks provider fallback defaults.
702 Optionally a raw, stored value can be retrieved by setting skip_interpolation to True.
703 This is useful for example when we want to write the default value to a file, and we don't
704 want the interpolation to happen as it is going to be done later when the config is read.
705
706 :param section: section of the config
707 :param key: key to use
708 :param fallback: fallback value to use
709 :param raw: if raw, then interpolation will be reversed
710 :param kwargs: other args
711 :return:
712 """
713 value = self._default_values.get(section, key, fallback=VALUE_NOT_FOUND_SENTINEL, **kwargs)
714 # Provider metadata has higher priority than cfg fallback — check it first.
715 if value is VALUE_NOT_FOUND_SENTINEL and self._use_providers_configuration:
716 value = self._provider_metadata_config_fallback_default_values.get(
717 section, key, fallback=VALUE_NOT_FOUND_SENTINEL, **kwargs
718 )
719 if value is VALUE_NOT_FOUND_SENTINEL and self._use_providers_configuration:
720 value = self._provider_cfg_config_fallback_default_values.get(
721 section, key, fallback=VALUE_NOT_FOUND_SENTINEL, **kwargs
722 )
723 if value is VALUE_NOT_FOUND_SENTINEL:
724 value = fallback
725 if raw and isinstance(value, str):
726 return value.replace("%", "%%")
727 return value
728
729 def _get_custom_secret_backend(self, worker_mode: bool = False) -> Any | None:
730 """
731 Get Secret Backend if defined in airflow.cfg.
732
733 Conditionally selects the section, key and kwargs key based on whether it is called from worker or not.
734 """
735 section, key = CUSTOM_BACKEND_CONFIG_KEYS["worker" if worker_mode else "general"]
736 kwargs_key = "secrets_backend_kwargs" if worker_mode else "backend_kwargs"
737
738 secrets_backend_cls = self.getimport(section=section, key=key)
739
740 if not secrets_backend_cls:
741 if worker_mode:
742 # if we find no secrets backend for worker, return that of secrets backend
743 section, key = CUSTOM_BACKEND_CONFIG_KEYS["general"]
744 secrets_backend_cls = self.getimport(section=section, key=key)
745 if not secrets_backend_cls:
746 return None
747 # When falling back to secrets backend, use its kwargs
748 kwargs_key = "backend_kwargs"
749 else:
750 return None
751
752 try:
753 backend_kwargs = self.getjson(section=section, key=kwargs_key)
754 if not backend_kwargs:
755 backend_kwargs = {}
756 elif not isinstance(backend_kwargs, dict):
757 raise ValueError("not a dict")
758 except AirflowConfigException:
759 log.warning("Failed to parse [%s] %s as JSON, defaulting to no kwargs.", section, kwargs_key)
760 backend_kwargs = {}
761 except ValueError:
762 log.warning("Failed to parse [%s] %s into a dict, defaulting to no kwargs.", section, kwargs_key)
763 backend_kwargs = {}
764
765 # Collect per-key overrides; they take precedence over the JSON blob.
766 env_prefix = _build_kwarg_env_prefix(section, kwargs_key)
767 backend_kwargs.update(_collect_kwarg_env_vars(env_prefix))
768
769 return secrets_backend_cls(**backend_kwargs)
770
771 def _get_config_value_from_secret_backend(self, config_key: str) -> str | None:
772 """
773 Get Config option values from Secret Backend.
774
775 Called by the shared parser's _get_secret_option() method as part of the lookup chain.
776 Uses _get_custom_secret_backend() to get the backend instance.
777
778 :param config_key: the config key to retrieve
779 :return: config value or None
780 """
781 try:
782 secrets_client = self._get_custom_secret_backend()
783 if not secrets_client:
784 return None
785 return secrets_client.get_config(config_key)
786 except Exception as e:
787 raise AirflowConfigException(
788 "Cannot retrieve config from alternative secrets backend. "
789 "Make sure it is configured properly and that the Backend "
790 "is accessible.\n"
791 f"{e}"
792 )
793
794 def _get_cmd_option_from_config_sources(
795 self, config_sources: ConfigSourcesType, section: str, key: str
796 ) -> str | None:
797 fallback_key = key + "_cmd"
798 if (section, key) in self.sensitive_config_values:
799 section_dict = config_sources.get(section)
800 if section_dict is not None:
801 command_value = section_dict.get(fallback_key)
802 if command_value is not None:
803 if isinstance(command_value, str):
804 command = command_value
805 else:
806 command = command_value[0]
807 return run_command(command)
808 return None
809
810 def _get_secret_option_from_config_sources(
811 self, config_sources: ConfigSourcesType, section: str, key: str
812 ) -> str | None:
813 fallback_key = key + "_secret"
814 if (section, key) in self.sensitive_config_values:
815 section_dict = config_sources.get(section)
816 if section_dict is not None:
817 secrets_path_value = section_dict.get(fallback_key)
818 if secrets_path_value is not None:
819 if isinstance(secrets_path_value, str):
820 secrets_path = secrets_path_value
821 else:
822 secrets_path = secrets_path_value[0]
823 return self._get_config_value_from_secret_backend(secrets_path)
824 return None
825
826 def _include_secrets(
827 self,
828 config_sources: ConfigSourcesType,
829 display_sensitive: bool,
830 display_source: bool,
831 raw: bool,
832 ):
833 for section, key in self.sensitive_config_values:
834 value: str | None = self._get_secret_option_from_config_sources(config_sources, section, key)
835 if value:
836 if not display_sensitive:
837 value = "< hidden >"
838 if display_source:
839 opt: str | tuple[str, str] = (value, "secret")
840 elif raw:
841 opt = value.replace("%", "%%")
842 else:
843 opt = value
844 config_sources.setdefault(section, {}).update({key: opt})
845 del config_sources[section][key + "_secret"]
846
847 def _include_commands(
848 self,
849 config_sources: ConfigSourcesType,
850 display_sensitive: bool,
851 display_source: bool,
852 raw: bool,
853 ):
854 for section, key in self.sensitive_config_values:
855 opt = self._get_cmd_option_from_config_sources(config_sources, section, key)
856 if not opt:
857 continue
858 opt_to_set: str | tuple[str, str] | None = opt
859 if not display_sensitive:
860 opt_to_set = "< hidden >"
861 if display_source:
862 opt_to_set = (str(opt_to_set), "cmd")
863 elif raw:
864 opt_to_set = str(opt_to_set).replace("%", "%%")
865 if opt_to_set is not None:
866 dict_to_update: dict[str, str | tuple[str, str]] = {key: opt_to_set}
867 config_sources.setdefault(section, {}).update(dict_to_update)
868 del config_sources[section][key + "_cmd"]
869
870 def _include_envs(
871 self,
872 config_sources: ConfigSourcesType,
873 display_sensitive: bool,
874 display_source: bool,
875 raw: bool,
876 ):
877 for env_var in [
878 os_environment for os_environment in os.environ if os_environment.startswith(ENV_VAR_PREFIX)
879 ]:
880 try:
881 _, section, key = env_var.split("__", 2)
882 opt = self._get_env_var_option(section, key)
883 except ValueError:
884 continue
885 if opt is None:
886 log.warning("Ignoring unknown env var '%s'", env_var)
887 continue
888 if not display_sensitive and env_var != self._env_var_name("core", "unit_test_mode"):
889 if self._names_sensitive_team_env_var(env_var):
890 # Covers the cmd/secret variants too; see is_sensitive_option.
891 opt = "< hidden >"
892 # Don't hide cmd/secret values here
893 elif not env_var.lower().endswith(("cmd", "secret")):
894 if (section, key) in self.sensitive_config_values:
895 opt = "< hidden >"
896 elif raw:
897 opt = opt.replace("%", "%%")
898 if display_source:
899 opt = (opt, "env var")
900
901 section = section.lower()
902 key = key.lower()
903 config_sources.setdefault(section, {}).update({key: opt})
904
905 def _filter_by_source(
906 self,
907 config_sources: ConfigSourcesType,
908 display_source: bool,
909 getter_func,
910 ):
911 """
912 Delete default configs from current configuration.
913
914 An OrderedDict of OrderedDicts, if it would conflict with special sensitive_config_values.
915
916 This is necessary because bare configs take precedence over the command
917 or secret key equivalents so if the current running config is
918 materialized with Airflow defaults they in turn override user set
919 command or secret key configs.
920
921 :param config_sources: The current configuration to operate on
922 :param display_source: If False, configuration options contain raw
923 values. If True, options are a tuple of (option_value, source).
924 Source is either 'airflow.cfg', 'default', 'env var', or 'cmd'.
925 :param getter_func: A callback function that gets the user configured
926 override value for a particular sensitive_config_values config.
927 :return: None, the given config_sources is filtered if necessary,
928 otherwise untouched.
929 """
930 for section, key in self.sensitive_config_values:
931 # Don't bother if we don't have section / key
932 if section not in config_sources or key not in config_sources[section]:
933 continue
934 # Check that there is something to override defaults
935 try:
936 getter_opt = getter_func(section, key)
937 except ValueError:
938 continue
939 if not getter_opt:
940 continue
941 # Check to see that there is a default value
942 if self.get_default_value(section, key) is None:
943 continue
944 # Check to see if bare setting is the same as defaults
945 if display_source:
946 # when display_source = true, we know that the config_sources contains tuple
947 opt, source = config_sources[section][key] # type: ignore
948 else:
949 opt = config_sources[section][key] # type: ignore[assignment]
950 if opt == self.get_default_value(section, key):
951 del config_sources[section][key]
952
953 @staticmethod
954 def _deprecated_value_is_set_in_config(
955 deprecated_section: str,
956 deprecated_key: str,
957 configs: Iterable[tuple[str, ConfigParser]],
958 ) -> bool:
959 for config_type, config in configs:
960 if config_type != "default":
961 with contextlib.suppress(NoSectionError):
962 deprecated_section_array = config.items(section=deprecated_section, raw=True)
963 if any(key == deprecated_key for key, _ in deprecated_section_array):
964 return True
965 return False
966
967 @staticmethod
968 def _deprecated_variable_is_set(deprecated_section: str, deprecated_key: str) -> bool:
969 return (
970 os.environ.get(f"{ENV_VAR_PREFIX}{deprecated_section.upper()}__{deprecated_key.upper()}")
971 is not None
972 )
973
974 @staticmethod
975 def _deprecated_command_is_set_in_config(
976 deprecated_section: str,
977 deprecated_key: str,
978 configs: Iterable[tuple[str, ConfigParser]],
979 ) -> bool:
980 return AirflowConfigParser._deprecated_value_is_set_in_config(
981 deprecated_section=deprecated_section, deprecated_key=deprecated_key + "_cmd", configs=configs
982 )
983
984 @staticmethod
985 def _deprecated_variable_command_is_set(deprecated_section: str, deprecated_key: str) -> bool:
986 return (
987 os.environ.get(f"{ENV_VAR_PREFIX}{deprecated_section.upper()}__{deprecated_key.upper()}_CMD")
988 is not None
989 )
990
991 @staticmethod
992 def _deprecated_secret_is_set_in_config(
993 deprecated_section: str,
994 deprecated_key: str,
995 configs: Iterable[tuple[str, ConfigParser]],
996 ) -> bool:
997 return AirflowConfigParser._deprecated_value_is_set_in_config(
998 deprecated_section=deprecated_section, deprecated_key=deprecated_key + "_secret", configs=configs
999 )
1000
1001 @staticmethod
1002 def _deprecated_variable_secret_is_set(deprecated_section: str, deprecated_key: str) -> bool:
1003 return (
1004 os.environ.get(f"{ENV_VAR_PREFIX}{deprecated_section.upper()}__{deprecated_key.upper()}_SECRET")
1005 is not None
1006 )
1007
1008 @staticmethod
1009 def _replace_config_with_display_sources(
1010 config_sources: ConfigSourcesType,
1011 configs: Iterable[tuple[str, ConfigParser]],
1012 configuration_description: dict[str, dict[str, Any]],
1013 display_source: bool,
1014 raw: bool,
1015 deprecated_options: dict[tuple[str, str], tuple[str, str, str]],
1016 include_env: bool,
1017 include_cmds: bool,
1018 include_secret: bool,
1019 ):
1020 for source_name, config in configs:
1021 sections = config.sections()
1022 for section in sections:
1023 AirflowConfigParser._replace_section_config_with_display_sources(
1024 config,
1025 config_sources,
1026 configuration_description,
1027 display_source,
1028 raw,
1029 section,
1030 source_name,
1031 deprecated_options,
1032 configs,
1033 include_env=include_env,
1034 include_cmds=include_cmds,
1035 include_secret=include_secret,
1036 )
1037
1038 @staticmethod
1039 def _replace_section_config_with_display_sources(
1040 config: ConfigParser,
1041 config_sources: ConfigSourcesType,
1042 configuration_description: dict[str, dict[str, Any]],
1043 display_source: bool,
1044 raw: bool,
1045 section: str,
1046 source_name: str,
1047 deprecated_options: dict[tuple[str, str], tuple[str, str, str]],
1048 configs: Iterable[tuple[str, ConfigParser]],
1049 include_env: bool,
1050 include_cmds: bool,
1051 include_secret: bool,
1052 ):
1053 sect = config_sources.setdefault(section, {})
1054 if isinstance(config, AirflowConfigParser):
1055 with config.suppress_future_warnings():
1056 items: Iterable[tuple[str, Any]] = config.items(section=section, raw=raw)
1057 else:
1058 items = config.items(section=section, raw=raw)
1059 for k, val in items:
1060 deprecated_section, deprecated_key, _ = deprecated_options.get((section, k), (None, None, None))
1061 if deprecated_section and deprecated_key:
1062 if source_name == "default":
1063 # If deprecated entry has some non-default value set for any of the sources requested,
1064 # We should NOT set default for the new entry (because it will override anything
1065 # coming from the deprecated ones)
1066 if AirflowConfigParser._deprecated_value_is_set_in_config(
1067 deprecated_section, deprecated_key, configs
1068 ):
1069 continue
1070 if include_env and AirflowConfigParser._deprecated_variable_is_set(
1071 deprecated_section, deprecated_key
1072 ):
1073 continue
1074 if include_cmds and (
1075 AirflowConfigParser._deprecated_variable_command_is_set(
1076 deprecated_section, deprecated_key
1077 )
1078 or AirflowConfigParser._deprecated_command_is_set_in_config(
1079 deprecated_section, deprecated_key, configs
1080 )
1081 ):
1082 continue
1083 if include_secret and (
1084 AirflowConfigParser._deprecated_variable_secret_is_set(
1085 deprecated_section, deprecated_key
1086 )
1087 or AirflowConfigParser._deprecated_secret_is_set_in_config(
1088 deprecated_section, deprecated_key, configs
1089 )
1090 ):
1091 continue
1092 if display_source:
1093 updated_source_name = source_name
1094 if source_name == "default":
1095 # defaults can come from other sources (default-<PROVIDER>) that should be used here
1096 source_description_section = configuration_description.get(section, {})
1097 source_description_key = source_description_section.get("options", {}).get(k, {})
1098 if source_description_key is not None:
1099 updated_source_name = source_description_key.get("source", source_name)
1100 sect[k] = (val, updated_source_name)
1101 else:
1102 sect[k] = val
1103
1104 def _warn_deprecate(
1105 self, section: str, key: str, deprecated_section: str, deprecated_name: str, extra_stacklevel: int
1106 ):
1107 """Warn about deprecated config option usage."""
1108 if section == deprecated_section:
1109 warnings.warn(
1110 f"The {deprecated_name} option in [{section}] has been renamed to {key} - "
1111 f"the old setting has been used, but please update your config.",
1112 DeprecationWarning,
1113 stacklevel=4 + extra_stacklevel,
1114 )
1115 else:
1116 warnings.warn(
1117 f"The {deprecated_name} option in [{deprecated_section}] has been moved to the {key} option "
1118 f"in [{section}] - the old setting has been used, but please update your config.",
1119 DeprecationWarning,
1120 stacklevel=4 + extra_stacklevel,
1121 )
1122
1123 @contextmanager
1124 def suppress_future_warnings(self):
1125 """
1126 Context manager to temporarily suppress future warnings.
1127
1128 This is a stub used by the shared parser's lookup methods when checking deprecated options.
1129 Subclasses can override this to customize warning suppression behavior.
1130
1131 :return: context manager that suppresses future warnings
1132 """
1133 suppress_future_warnings = self._suppress_future_warnings
1134 self._suppress_future_warnings = True
1135 yield self
1136 self._suppress_future_warnings = suppress_future_warnings
1137
1138 def _env_var_name(self, section: str, key: str, team_name: str | None = None) -> str:
1139 """Generate environment variable name for a config option."""
1140 team_component: str = f"{team_name.upper()}___" if team_name else ""
1141 return f"{ENV_VAR_PREFIX}{team_component}{section.replace('.', '_').upper()}__{key.upper()}"
1142
1143 def _get_env_var_option(self, section: str, key: str, team_name: str | None = None):
1144 """Get config option from environment variable."""
1145 env_var: str = self._env_var_name(section, key, team_name=team_name)
1146 if env_var in os.environ:
1147 return expand_env_var(os.environ[env_var])
1148 # alternatively AIRFLOW__{SECTION}__{KEY}_CMD (for a command)
1149 env_var_cmd = env_var + "_CMD"
1150 if env_var_cmd in os.environ:
1151 # if this is a valid command key...
1152 if (section, key) in self.sensitive_config_values:
1153 return run_command(os.environ[env_var_cmd])
1154 # alternatively AIRFLOW__{SECTION}__{KEY}_SECRET (to get from Secrets Backend)
1155 env_var_secret_path = env_var + "_SECRET"
1156 if env_var_secret_path in os.environ:
1157 # if this is a valid secret path...
1158 if (section, key) in self.sensitive_config_values:
1159 return self._get_config_value_from_secret_backend(os.environ[env_var_secret_path])
1160 return None
1161
1162 def _get_cmd_option(self, section: str, key: str):
1163 """Get config option from command execution."""
1164 fallback_key = key + "_cmd"
1165 if (section, key) in self.sensitive_config_values:
1166 if super().has_option(section, fallback_key):
1167 command = super().get(section, fallback_key)
1168 try:
1169 cmd_output = run_command(command)
1170 except AirflowConfigException as e:
1171 raise e
1172 except Exception as e:
1173 raise AirflowConfigException(
1174 f"Cannot run the command for the config section [{section}]{fallback_key}_cmd."
1175 f" Please check the {fallback_key} value."
1176 ) from e
1177 return cmd_output
1178 return None
1179
1180 def _get_secret_option(self, section: str, key: str) -> str | None:
1181 """Get Config option values from Secret Backend."""
1182 fallback_key = key + "_secret"
1183 if (section, key) in self.sensitive_config_values:
1184 if super().has_option(section, fallback_key):
1185 secrets_path = super().get(section, fallback_key)
1186 return self._get_config_value_from_secret_backend(secrets_path)
1187 return None
1188
1189 def _get_environment_variables(
1190 self,
1191 deprecated_key: str | None,
1192 deprecated_section: str | None,
1193 key: str,
1194 section: str,
1195 issue_warning: bool = True,
1196 extra_stacklevel: int = 0,
1197 **kwargs,
1198 ) -> str | ValueNotFound:
1199 """Get config option from environment variables."""
1200 team_name = kwargs.get("team_name", None)
1201 option = self._get_env_var_option(section, key, team_name=team_name)
1202 if option is not None:
1203 return option
1204 if deprecated_section and deprecated_key:
1205 with self.suppress_future_warnings():
1206 option = self._get_env_var_option(deprecated_section, deprecated_key, team_name=team_name)
1207 if option is not None:
1208 if issue_warning:
1209 self._warn_deprecate(section, key, deprecated_section, deprecated_key, extra_stacklevel)
1210 return option
1211 return VALUE_NOT_FOUND_SENTINEL
1212
1213 def _get_option_from_config_file(
1214 self,
1215 deprecated_key: str | None,
1216 deprecated_section: str | None,
1217 key: str,
1218 section: str,
1219 issue_warning: bool = True,
1220 extra_stacklevel: int = 0,
1221 **kwargs,
1222 ) -> str | ValueNotFound:
1223 """Get config option from config file."""
1224 if team_name := kwargs.get("team_name", None):
1225 section = team_section_name(team_name, section)
1226 # since this is the last lookup that supports team_name, pop it
1227 kwargs.pop("team_name")
1228 if super().has_option(section, key):
1229 return expand_env_var(super().get(section, key, **kwargs))
1230 if deprecated_section and deprecated_key:
1231 if super().has_option(deprecated_section, deprecated_key):
1232 if issue_warning:
1233 self._warn_deprecate(section, key, deprecated_section, deprecated_key, extra_stacklevel)
1234 with self.suppress_future_warnings():
1235 return expand_env_var(super().get(deprecated_section, deprecated_key, **kwargs))
1236 return VALUE_NOT_FOUND_SENTINEL
1237
1238 def _get_option_from_commands(
1239 self,
1240 deprecated_key: str | None,
1241 deprecated_section: str | None,
1242 key: str,
1243 section: str,
1244 issue_warning: bool = True,
1245 extra_stacklevel: int = 0,
1246 **kwargs,
1247 ) -> str | ValueNotFound:
1248 """Get config option from command execution."""
1249 if kwargs.get("team_name", None):
1250 # Commands based team config fetching is not currently supported
1251 return VALUE_NOT_FOUND_SENTINEL
1252 option = self._get_cmd_option(section, key)
1253 if option:
1254 return option
1255 if deprecated_section and deprecated_key:
1256 with self.suppress_future_warnings():
1257 option = self._get_cmd_option(deprecated_section, deprecated_key)
1258 if option:
1259 if issue_warning:
1260 self._warn_deprecate(section, key, deprecated_section, deprecated_key, extra_stacklevel)
1261 return option
1262 return VALUE_NOT_FOUND_SENTINEL
1263
1264 def _get_option_from_secrets(
1265 self,
1266 deprecated_key: str | None,
1267 deprecated_section: str | None,
1268 key: str,
1269 section: str,
1270 issue_warning: bool = True,
1271 extra_stacklevel: int = 0,
1272 **kwargs,
1273 ) -> str | ValueNotFound:
1274 """Get config option from secrets backend."""
1275 if kwargs.get("team_name", None):
1276 # Secrets based team config fetching is not currently supported
1277 return VALUE_NOT_FOUND_SENTINEL
1278 option = self._get_secret_option(section, key)
1279 if option:
1280 return option
1281 if deprecated_section and deprecated_key:
1282 with self.suppress_future_warnings():
1283 option = self._get_secret_option(deprecated_section, deprecated_key)
1284 if option:
1285 if issue_warning:
1286 self._warn_deprecate(section, key, deprecated_section, deprecated_key, extra_stacklevel)
1287 return option
1288 return VALUE_NOT_FOUND_SENTINEL
1289
1290 def _get_option_from_defaults(
1291 self,
1292 deprecated_key: str | None,
1293 deprecated_section: str | None,
1294 key: str,
1295 section: str,
1296 issue_warning: bool = True,
1297 extra_stacklevel: int = 0,
1298 team_name: str | None = None,
1299 **kwargs,
1300 ) -> str | ValueNotFound:
1301 """Get config option from default values."""
1302 if self.get_default_value(section, key) is not None or "fallback" in kwargs:
1303 return expand_env_var(self.get_default_value(section, key, **kwargs))
1304 return VALUE_NOT_FOUND_SENTINEL
1305
1306 def _resolve_deprecated_lookup(
1307 self,
1308 section: str,
1309 key: str,
1310 lookup_from_deprecated: bool,
1311 extra_stacklevel: int = 0,
1312 ) -> tuple[str, str, str | None, str | None, bool]:
1313 """
1314 Resolve deprecated section/key mappings and determine deprecated values.
1315
1316 :param section: Section name (will be lowercased)
1317 :param key: Key name (will be lowercased)
1318 :param lookup_from_deprecated: Whether to lookup from deprecated options
1319 :param extra_stacklevel: Extra stack level for warnings
1320 :return: Tuple of (resolved_section, resolved_key, deprecated_section, deprecated_key, warning_emitted)
1321 """
1322 section = section.lower()
1323 key = key.lower()
1324 warning_emitted = False
1325 deprecated_section: str | None = None
1326 deprecated_key: str | None = None
1327
1328 if not lookup_from_deprecated:
1329 return section, key, deprecated_section, deprecated_key, warning_emitted
1330
1331 option_description = self.configuration_description.get(section, {}).get("options", {}).get(key, {})
1332 if option_description.get("deprecated"):
1333 deprecation_reason = option_description.get("deprecation_reason", "")
1334 warnings.warn(
1335 f"The '{key}' option in section {section} is deprecated. {deprecation_reason}",
1336 DeprecationWarning,
1337 stacklevel=2 + extra_stacklevel,
1338 )
1339 # For the cases in which we rename whole sections
1340 if section in self.inversed_deprecated_sections:
1341 deprecated_section, deprecated_key = (section, key)
1342 section = self.inversed_deprecated_sections[section]
1343 if not self._suppress_future_warnings:
1344 warnings.warn(
1345 f"The config section [{deprecated_section}] has been renamed to "
1346 f"[{section}]. Please update your `conf.get*` call to use the new name",
1347 FutureWarning,
1348 stacklevel=2 + extra_stacklevel,
1349 )
1350 # Don't warn about individual rename if the whole section is renamed
1351 warning_emitted = True
1352 elif (section, key) in self.inversed_deprecated_options:
1353 # Handle using deprecated section/key instead of the new section/key
1354 new_section, new_key = self.inversed_deprecated_options[(section, key)]
1355 if not self._suppress_future_warnings and not warning_emitted:
1356 warnings.warn(
1357 f"section/key [{section}/{key}] has been deprecated, you should use"
1358 f"[{new_section}/{new_key}] instead. Please update your `conf.get*` call to use the "
1359 "new name",
1360 FutureWarning,
1361 stacklevel=2 + extra_stacklevel,
1362 )
1363 warning_emitted = True
1364 deprecated_section, deprecated_key = section, key
1365 section, key = (new_section, new_key)
1366 elif section in self.deprecated_sections:
1367 # When accessing the new section name, make sure we check under the old config name
1368 deprecated_key = key
1369 deprecated_section = self.deprecated_sections[section][0]
1370 else:
1371 deprecated_section, deprecated_key, _ = self.deprecated_options.get(
1372 (section, key), (None, None, None)
1373 )
1374
1375 return section, key, deprecated_section, deprecated_key, warning_emitted
1376
1377 def load_providers_configuration(self) -> None:
1378 """
1379 Load configuration for providers.
1380
1381 .. deprecated:: 3.2.0
1382 Provider configuration is now loaded lazily via the ``configuration_description``
1383 cached property. This method is kept for backwards compatibility and will be
1384 removed in a future version.
1385 """
1386 warnings.warn(
1387 "load_providers_configuration() is deprecated. "
1388 "Provider configuration is now loaded lazily via the "
1389 "`configuration_description` cached property.",
1390 DeprecationWarning,
1391 stacklevel=2,
1392 )
1393 self._use_providers_configuration = True
1394 self._invalidate_provider_flag_caches()
1395
1396 def restore_core_default_configuration(self) -> None:
1397 """
1398 Restore the parser state before provider-contributed sections were loaded.
1399
1400 .. deprecated:: 3.2.0
1401 Use ``make_sure_configuration_loaded(with_providers=False)`` context manager
1402 instead. This method is kept for backwards compatibility and will be removed
1403 in a future version.
1404 """
1405 warnings.warn(
1406 "restore_core_default_configuration() is deprecated. "
1407 "Use `make_sure_configuration_loaded(with_providers=False)` instead.",
1408 DeprecationWarning,
1409 stacklevel=2,
1410 )
1411 self._use_providers_configuration = False
1412 self._invalidate_provider_flag_caches()
1413
1414 @overload # type: ignore[override]
1415 def get(self, section: str, key: str, fallback: str = ..., **kwargs) -> str: ...
1416
1417 @overload # type: ignore[override]
1418 def get(self, section: str, key: str, **kwargs) -> str | None: ...
1419
1420 def get( # type: ignore[misc, override]
1421 self,
1422 section: str,
1423 key: str,
1424 suppress_warnings: bool = False,
1425 lookup_from_deprecated: bool = True,
1426 _extra_stacklevel: int = 0,
1427 team_name: str | None = None,
1428 **kwargs,
1429 ) -> str | None:
1430 """
1431 Get config value by iterating through lookup sequence.
1432
1433 Priority order is defined by _lookup_sequence property.
1434 """
1435 section, key, deprecated_section, deprecated_key, warning_emitted = self._resolve_deprecated_lookup(
1436 section=section,
1437 key=key,
1438 lookup_from_deprecated=lookup_from_deprecated,
1439 extra_stacklevel=_extra_stacklevel,
1440 )
1441
1442 if team_name is not None:
1443 kwargs["team_name"] = team_name
1444
1445 for lookup_method in self._lookup_sequence:
1446 value = lookup_method(
1447 deprecated_key=deprecated_key,
1448 deprecated_section=deprecated_section,
1449 key=key,
1450 section=section,
1451 issue_warning=not warning_emitted,
1452 extra_stacklevel=_extra_stacklevel,
1453 **kwargs,
1454 )
1455 if value is not VALUE_NOT_FOUND_SENTINEL:
1456 return value
1457
1458 # Check if fallback was explicitly provided (even if None)
1459 if "fallback" in kwargs:
1460 return kwargs["fallback"]
1461
1462 if not suppress_warnings:
1463 log.warning("section/key [%s/%s] not found in config", section, key)
1464
1465 raise AirflowConfigException(f"section/key [{section}/{key}] not found in config")
1466
1467 def getboolean(self, section: str, key: str, **kwargs) -> bool: # type: ignore[override]
1468 """Get config value as boolean."""
1469 val = str(self.get(section, key, _extra_stacklevel=1, **kwargs)).lower().strip()
1470 if "#" in val:
1471 val = val.split("#")[0].strip()
1472 if val in ("t", "true", "1"):
1473 return True
1474 if val in ("f", "false", "0"):
1475 return False
1476 raise AirflowConfigException(
1477 f'Failed to convert value to bool. Please check "{key}" key in "{section}" section. '
1478 f'Current value: "{val}".'
1479 )
1480
1481 def getint(self, section: str, key: str, **kwargs) -> int: # type: ignore[override]
1482 """Get config value as integer."""
1483 val = self.get(section, key, _extra_stacklevel=1, **kwargs)
1484 if val is None:
1485 raise AirflowConfigException(
1486 f"Failed to convert value None to int. "
1487 f'Please check "{key}" key in "{section}" section is set.'
1488 )
1489 try:
1490 return int(val)
1491 except ValueError:
1492 try:
1493 if (float_val := float(val)) != (int_val := int(float_val)):
1494 raise ValueError
1495 return int_val
1496 except (ValueError, OverflowError):
1497 raise AirflowConfigException(
1498 f'Failed to convert value to int. Please check "{key}" key in "{section}" section. '
1499 f'Current value: "{val}".'
1500 )
1501
1502 def getfloat(self, section: str, key: str, **kwargs) -> float: # type: ignore[override]
1503 """Get config value as float."""
1504 val = self.get(section, key, _extra_stacklevel=1, **kwargs)
1505 if val is None:
1506 raise AirflowConfigException(
1507 f"Failed to convert value None to float. "
1508 f'Please check "{key}" key in "{section}" section is set.'
1509 )
1510 try:
1511 return float(val)
1512 except ValueError:
1513 raise AirflowConfigException(
1514 f'Failed to convert value to float. Please check "{key}" key in "{section}" section. '
1515 f'Current value: "{val}".'
1516 )
1517
1518 def getlist(self, section: str, key: str, delimiter=",", **kwargs):
1519 """Get config value as list."""
1520 val = self.get(section, key, **kwargs)
1521
1522 if isinstance(val, list) or val is None:
1523 # `get` will always return a (possibly-empty) string, so the only way we can
1524 # have these types is with `fallback=` was specified. So just return it.
1525 return val
1526
1527 if val == "":
1528 return []
1529
1530 try:
1531 return [item.strip() for item in val.split(delimiter)]
1532 except Exception:
1533 raise AirflowConfigException(
1534 f'Failed to parse value to a list. Please check "{key}" key in "{section}" section. '
1535 f'Current value: "{val}".'
1536 )
1537
1538 E = TypeVar("E", bound=Enum)
1539
1540 def getenum(self, section: str, key: str, enum_class: type[E], **kwargs) -> E:
1541 """Get config value as enum."""
1542 val = self.get(section, key, **kwargs)
1543 enum_names = [enum_item.name for enum_item in enum_class]
1544
1545 if val is None:
1546 raise AirflowConfigException(
1547 f'Failed to convert value. Please check "{key}" key in "{section}" section. '
1548 f'Current value: "{val}" and it must be one of {", ".join(enum_names)}'
1549 )
1550
1551 try:
1552 return enum_class[val]
1553 except KeyError:
1554 if "fallback" in kwargs and kwargs["fallback"] in enum_names:
1555 return enum_class[kwargs["fallback"]]
1556 raise AirflowConfigException(
1557 f'Failed to convert value. Please check "{key}" key in "{section}" section. '
1558 f"the value must be one of {', '.join(enum_names)}"
1559 )
1560
1561 def getenumlist(self, section: str, key: str, enum_class: type[E], delimiter=",", **kwargs) -> list[E]:
1562 """Get config value as list of enums."""
1563 kwargs.setdefault("fallback", [])
1564 string_list = self.getlist(section, key, delimiter, **kwargs)
1565
1566 enum_names = [enum_item.name for enum_item in enum_class]
1567 enum_list = []
1568
1569 for val in string_list:
1570 try:
1571 enum_list.append(enum_class[val])
1572 except KeyError:
1573 log.warning(
1574 "Failed to convert value %r. Please check %s key in %s section. "
1575 "it must be one of %s, if not the value is ignored",
1576 val,
1577 key,
1578 section,
1579 ", ".join(enum_names),
1580 )
1581
1582 return enum_list
1583
1584 def getimport(self, section: str, key: str, **kwargs) -> Any:
1585 """
1586 Read options, import the full qualified name, and return the object.
1587
1588 In case of failure, it throws an exception with the key and section names
1589
1590 :return: The object or None, if the option is empty
1591 """
1592 # Fixed: use self.get() instead of conf.get()
1593 full_qualified_path = self.get(section=section, key=key, **kwargs)
1594 if not full_qualified_path:
1595 return None
1596
1597 try:
1598 # Import here to avoid circular dependency
1599 from ..module_loading import import_string
1600
1601 return import_string(full_qualified_path)
1602 except ImportError as e:
1603 log.warning(e)
1604 raise AirflowConfigException(
1605 f'The object could not be loaded. Please check "{key}" key in "{section}" section. '
1606 f'Current value: "{full_qualified_path}".'
1607 )
1608
1609 def getjson(
1610 self, section: str, key: str, fallback=None, **kwargs
1611 ) -> dict | list | str | int | float | None:
1612 """
1613 Return a config value parsed from a JSON string.
1614
1615 ``fallback`` is *not* JSON parsed but used verbatim when no config value is given.
1616 """
1617 try:
1618 data = self.get(section=section, key=key, fallback=None, _extra_stacklevel=1, **kwargs)
1619 except (NoSectionError, NoOptionError):
1620 data = None
1621
1622 if data is None or data == "":
1623 return fallback
1624
1625 try:
1626 return json.loads(data)
1627 except JSONDecodeError as e:
1628 raise AirflowConfigException(f"Unable to parse [{section}] {key!r} as valid json") from e
1629
1630 def gettimedelta(
1631 self, section: str, key: str, fallback: Any = None, **kwargs
1632 ) -> datetime.timedelta | None:
1633 """
1634 Get the config value for the given section and key, and convert it into datetime.timedelta object.
1635
1636 If the key is missing, then it is considered as `None`.
1637
1638 :param section: the section from the config
1639 :param key: the key defined in the given section
1640 :param fallback: fallback value when no config value is given, defaults to None
1641 :raises AirflowConfigException: raised because ValueError or OverflowError
1642 :return: datetime.timedelta(seconds=<config_value>) or None
1643 """
1644 val = self.get(section, key, fallback=fallback, _extra_stacklevel=1, **kwargs)
1645
1646 if val:
1647 # the given value must be convertible to integer
1648 try:
1649 int_val = int(val)
1650 except ValueError:
1651 raise AirflowConfigException(
1652 f'Failed to convert value to int. Please check "{key}" key in "{section}" section. '
1653 f'Current value: "{val}".'
1654 )
1655
1656 try:
1657 return datetime.timedelta(seconds=int_val)
1658 except OverflowError as err:
1659 raise AirflowConfigException(
1660 f"Failed to convert value to timedelta in `seconds`. "
1661 f"{err}. "
1662 f'Please check "{key}" key in "{section}" section. Current value: "{val}".'
1663 )
1664
1665 return fallback
1666
1667 def get_mandatory_value(self, section: str, key: str, **kwargs) -> str:
1668 """Get mandatory config value, raising ValueError if not found."""
1669 value = self.get(section, key, _extra_stacklevel=1, **kwargs)
1670 if value is None:
1671 raise ValueError(f"The value {section}/{key} should be set!")
1672 return value
1673
1674 def get_mandatory_list_value(self, section: str, key: str, **kwargs) -> list[str]:
1675 """Get mandatory config value as list, raising ValueError if not found."""
1676 value = self.getlist(section, key, **kwargs)
1677 if value is None:
1678 raise ValueError(f"The value {section}/{key} should be set!")
1679 return value
1680
1681 def read( # type: ignore[override]
1682 self,
1683 filenames: str | bytes | os.PathLike | Iterable[str | bytes | os.PathLike],
1684 encoding: str | None = None,
1685 ) -> list[str]:
1686 return super().read(filenames=filenames, encoding=encoding) # type: ignore[arg-type,return-value]
1687
1688 def read_dict( # type: ignore[override]
1689 self, dictionary: dict[str, dict[str, Any]], source: str = "<dict>"
1690 ) -> None:
1691 """
1692 We define a different signature here to add better type hints and checking.
1693
1694 :param dictionary: dictionary to read from
1695 :param source: source to be used to store the configuration
1696 :return:
1697 """
1698 super().read_dict(dictionary=dictionary, source=source)
1699
1700 def _has_section_in_any_defaults(self, section: str) -> bool:
1701 """Check if section exists in core defaults or provider fallback defaults."""
1702 if self._default_values.has_section(section):
1703 return True
1704 if self._use_providers_configuration:
1705 if self._provider_cfg_config_fallback_default_values.has_section(section):
1706 return True
1707 if self._provider_metadata_config_fallback_default_values.has_section(section):
1708 return True
1709 return False
1710
1711 def get_sections_including_defaults(self) -> list[str]:
1712 """
1713 Retrieve all sections from the configuration parser, including sections defined by built-in defaults.
1714
1715 :return: list of section names
1716 """
1717 sections_from_config = self.sections()
1718 sections_from_description = list(self.configuration_description.keys())
1719 return list(dict.fromkeys(itertools.chain(sections_from_description, sections_from_config)))
1720
1721 def get_options_including_defaults(self, section: str) -> list[str]:
1722 """
1723 Retrieve all possible options from the configuration parser for the section given.
1724
1725 Includes options defined by built-in defaults.
1726
1727 :param section: section name
1728 :return: list of option names for the section given
1729 """
1730 my_own_options = self.options(section) if self.has_section(section) else []
1731 all_options_from_defaults = list(
1732 self.configuration_description.get(section, {}).get("options", {}).keys()
1733 )
1734 return list(dict.fromkeys(itertools.chain(all_options_from_defaults, my_own_options)))
1735
1736 def has_option( # type: ignore[override]
1737 self, section: str, option: str, lookup_from_deprecated: bool = True, **kwargs
1738 ) -> bool:
1739 """
1740 Check if option is defined.
1741
1742 Uses self.get() to avoid reimplementing the priority order of config variables
1743 (env, config, cmd, defaults).
1744
1745 :param section: section to get option from
1746 :param option: option to get
1747 :param lookup_from_deprecated: If True, check if the option is defined in deprecated sections
1748 :param kwargs: additional keyword arguments to pass to get(), such as team_name
1749 :return:
1750 """
1751 try:
1752 value = self.get(
1753 section,
1754 option,
1755 fallback=VALUE_NOT_FOUND_SENTINEL,
1756 _extra_stacklevel=1,
1757 suppress_warnings=True,
1758 lookup_from_deprecated=lookup_from_deprecated,
1759 **kwargs,
1760 )
1761 if value is VALUE_NOT_FOUND_SENTINEL:
1762 return False
1763 return True
1764 except (NoOptionError, NoSectionError, AirflowConfigException):
1765 return False
1766
1767 def set(self, section: str, option: str, value: str | None = None) -> None: # type: ignore[override]
1768 """
1769 Set an option to the given value.
1770
1771 This override just makes sure the section and option are lower case, to match what we do in `get`.
1772 """
1773 section = section.lower()
1774 option = option.lower()
1775 defaults = self.configuration_description or {}
1776 if not self.has_section(section) and section in defaults:
1777 # Trying to set a key in a section that exists in default, but not in the user config;
1778 # automatically create it
1779 self.add_section(section)
1780 super().set(section, option, value)
1781
1782 def remove_option(self, section: str, option: str, remove_default: bool = True): # type: ignore[override]
1783 """
1784 Remove an option if it exists in config from a file or default config.
1785
1786 If both of config have the same option, this removes the option
1787 in both configs unless remove_default=False.
1788 """
1789 section = section.lower()
1790 option = option.lower()
1791 if super().has_option(section, option):
1792 super().remove_option(section, option)
1793
1794 if remove_default and self._default_values.has_option(section, option):
1795 self._default_values.remove_option(section, option)
1796
1797 def optionxform(self, optionstr: str) -> str:
1798 """
1799 Transform option names on every read, get, or set operation.
1800
1801 This changes from the default behaviour of ConfigParser from lower-casing
1802 to instead be case-preserving.
1803
1804 :param optionstr:
1805 :return:
1806 """
1807 return optionstr
1808
1809 def as_dict(
1810 self,
1811 display_source: bool = False,
1812 display_sensitive: bool = False,
1813 raw: bool = False,
1814 include_env: bool = True,
1815 include_cmds: bool = True,
1816 include_secret: bool = True,
1817 ) -> ConfigSourcesType:
1818 """
1819 Return the current configuration as an OrderedDict of OrderedDicts.
1820
1821 When materializing current configuration Airflow defaults are
1822 materialized along with user set configs. If any of the `include_*`
1823 options are False then the result of calling command or secret key
1824 configs do not override Airflow defaults and instead are passed through.
1825 In order to then avoid Airflow defaults from overwriting user set
1826 command or secret key configs we filter out bare sensitive_config_values
1827 that are set to Airflow defaults when command or secret key configs
1828 produce different values.
1829
1830 :param display_source: If False, the option value is returned. If True,
1831 a tuple of (option_value, source) is returned. Source is either
1832 'airflow.cfg', 'default', 'env var', or 'cmd'.
1833 :param display_sensitive: If True, the values of options set by env
1834 vars and bash commands will be displayed. If False, those options
1835 are shown as '< hidden >'
1836 :param raw: Should the values be output as interpolated values, or the
1837 "raw" form that can be fed back in to ConfigParser
1838 :param include_env: Should the value of configuration from AIRFLOW__
1839 environment variables be included or not
1840 :param include_cmds: Should the result of calling any ``*_cmd`` config be
1841 set (True, default), or should the _cmd options be left as the
1842 command to run (False)
1843 :param include_secret: Should the result of calling any ``*_secret`` config be
1844 set (True, default), or should the _secret options be left as the
1845 path to get the secret from (False)
1846 :return: Dictionary, where the key is the name of the section and the content is
1847 the dictionary with the name of the parameter and its value.
1848 """
1849 if not display_sensitive:
1850 # We want to hide the sensitive values at the appropriate methods
1851 # since envs from cmds, secrets can be read at _include_envs method
1852 if not all([include_env, include_cmds, include_secret]):
1853 raise ValueError(
1854 "If display_sensitive is false, then include_env, "
1855 "include_cmds, include_secret must all be set as True"
1856 )
1857
1858 config_sources: ConfigSourcesType = {}
1859
1860 # We check sequentially all those sources and the last one we saw it in will "win"
1861 configs = self._config_sources_for_as_dict
1862
1863 self._replace_config_with_display_sources(
1864 config_sources,
1865 configs,
1866 self.configuration_description,
1867 display_source,
1868 raw,
1869 self.deprecated_options,
1870 include_cmds=include_cmds,
1871 include_env=include_env,
1872 include_secret=include_secret,
1873 )
1874
1875 # add env vars and overwrite because they have priority
1876 if include_env:
1877 self._include_envs(config_sources, display_sensitive, display_source, raw)
1878 else:
1879 self._filter_by_source(config_sources, display_source, self._get_env_var_option)
1880
1881 # add bash commands
1882 if include_cmds:
1883 self._include_commands(config_sources, display_sensitive, display_source, raw)
1884 else:
1885 self._filter_by_source(config_sources, display_source, self._get_cmd_option)
1886
1887 # add config from secret backends
1888 if include_secret:
1889 self._include_secrets(config_sources, display_sensitive, display_source, raw)
1890 else:
1891 self._filter_by_source(config_sources, display_source, self._get_secret_option)
1892
1893 if not display_sensitive:
1894 # This ensures the ones from config file is hidden too
1895 # if they are not provided through env, cmd and secret
1896 # The collected options are walked (rather than the registered sensitive ones) so that
1897 # team scoped sections are covered as well - they are named after the team, not after
1898 # the base section the option is registered under.
1899 hidden = "< hidden >"
1900 for section, options in config_sources.items():
1901 for key, value in list(options.items()):
1902 if not value or not self.is_sensitive_option(section, key):
1903 continue
1904 if display_source:
1905 source = value[1]
1906 options[key] = (hidden, source)
1907 else:
1908 options[key] = hidden
1909
1910 return config_sources
1911
1912 def _write_option_header(
1913 self,
1914 file: IO[str],
1915 option: str,
1916 extra_spacing: bool,
1917 include_descriptions: bool,
1918 include_env_vars: bool,
1919 include_examples: bool,
1920 include_sources: bool,
1921 section_config_description: dict[str, dict[str, Any]],
1922 section_to_write: str,
1923 sources_dict: ConfigSourcesType,
1924 ) -> tuple[bool, bool]:
1925 """
1926 Write header for configuration option.
1927
1928 Returns tuple of (should_continue, needs_separation) where needs_separation should be
1929 set if the option needs additional separation to visually separate it from the next option.
1930 """
1931 option_config_description = (
1932 section_config_description.get("options", {}).get(option, {})
1933 if section_config_description
1934 else {}
1935 )
1936 description = option_config_description.get("description")
1937 needs_separation = False
1938 if description and include_descriptions:
1939 for line in description.splitlines():
1940 file.write(f"# {line}\n")
1941 needs_separation = True
1942 example = option_config_description.get("example")
1943 if example is not None and include_examples:
1944 if extra_spacing:
1945 file.write("#\n")
1946 example_lines = example.splitlines()
1947 example = "\n# ".join(example_lines)
1948 file.write(f"# Example: {option} = {example}\n")
1949 needs_separation = True
1950 if include_sources and sources_dict:
1951 sources_section = sources_dict.get(section_to_write)
1952 value_with_source = sources_section.get(option) if sources_section else None
1953 if value_with_source is None:
1954 file.write("#\n# Source: not defined\n")
1955 else:
1956 file.write(f"#\n# Source: {value_with_source[1]}\n")
1957 needs_separation = True
1958 if include_env_vars:
1959 file.write(f"#\n# Variable: AIRFLOW__{section_to_write.upper()}__{option.upper()}\n")
1960 if extra_spacing:
1961 file.write("#\n")
1962 needs_separation = True
1963 return True, needs_separation
1964
1965 def is_template(self, section: str, key) -> bool:
1966 """
1967 Return whether the value is templated.
1968
1969 :param section: section of the config
1970 :param key: key in the section
1971 :return: True if the value is templated
1972 """
1973 return _is_template(self.configuration_description, section, key)
1974
1975 def getsection(self, section: str, team_name: str | None = None) -> ConfigOptionsDictType | None:
1976 """
1977 Return the section as a dict.
1978
1979 Values are converted to int, float, bool as required.
1980
1981 :param section: section from the config
1982 :param team_name: optional team name for team-specific configuration lookup
1983 """
1984 # Handle team-specific section lookup for config file
1985 config_section = team_section_name(team_name, section) if team_name else section
1986
1987 if not self.has_section(config_section) and not self._has_section_in_any_defaults(config_section):
1988 return None
1989 if self._default_values.has_section(config_section):
1990 _section: ConfigOptionsDictType = dict(self._default_values.items(config_section))
1991 else:
1992 _section = {}
1993
1994 if self.has_section(config_section):
1995 _section.update(self.items(config_section))
1996
1997 # Use section (not config_section) for env var lookup - team_name is handled by _env_var_name
1998 section_prefix = self._env_var_name(section, "", team_name=team_name)
1999 for env_var in sorted(os.environ.keys()):
2000 if env_var.startswith(section_prefix):
2001 key = env_var.replace(section_prefix, "")
2002 if key.endswith("_CMD"):
2003 key = key[:-4]
2004 key = key.lower()
2005 _section[key] = self._get_env_var_option(section, key, team_name=team_name)
2006
2007 for key, val in _section.items():
2008 if val is None:
2009 raise AirflowConfigException(
2010 f"Failed to convert value automatically. "
2011 f'Please check "{key}" key in "{section}" section is set.'
2012 )
2013 try:
2014 _section[key] = int(val)
2015 except ValueError:
2016 try:
2017 _section[key] = float(val)
2018 except ValueError:
2019 if isinstance(val, str) and val.lower() in ("t", "true"):
2020 _section[key] = True
2021 elif isinstance(val, str) and val.lower() in ("f", "false"):
2022 _section[key] = False
2023 return _section
2024
2025 @staticmethod
2026 def _write_section_header(
2027 file: IO[str],
2028 include_descriptions: bool,
2029 section_config_description: dict[str, str],
2030 section_to_write: str,
2031 ) -> None:
2032 """Write header for configuration section."""
2033 file.write(f"[{section_to_write}]\n")
2034 section_description = section_config_description.get("description")
2035 if section_description and include_descriptions:
2036 for line in section_description.splitlines():
2037 file.write(f"# {line}\n")
2038 file.write("\n")
2039
2040 def _write_value(
2041 self,
2042 file: IO[str],
2043 option: str,
2044 comment_out_everything: bool,
2045 needs_separation: bool,
2046 only_defaults: bool,
2047 section_to_write: str,
2048 hide_sensitive: bool,
2049 is_sensitive: bool,
2050 show_values: bool = False,
2051 ):
2052 default_value = self.get_default_value(section_to_write, option, raw=True)
2053 if only_defaults:
2054 value = default_value
2055 else:
2056 value = self.get(section_to_write, option, fallback=default_value, raw=True)
2057 if not show_values:
2058 file.write(f"# {option} = \n")
2059 else:
2060 if hide_sensitive and is_sensitive:
2061 value = "< hidden >"
2062 else:
2063 pass
2064 if value is None:
2065 file.write(f"# {option} = \n")
2066 else:
2067 if comment_out_everything:
2068 value_lines = value.splitlines()
2069 value = "\n# ".join(value_lines)
2070 file.write(f"# {option} = {value}\n")
2071 else:
2072 if "\n" in value:
2073 try:
2074 value = json.dumps(json.loads(value), indent=4)
2075 value = value.replace(
2076 "\n", "\n "
2077 ) # indent multi-line JSON to satisfy configparser format
2078 except JSONDecodeError:
2079 pass
2080 file.write(f"{option} = {value}\n")
2081 if needs_separation:
2082 file.write("\n")
2083
2084 def write( # type: ignore[override]
2085 self,
2086 file: IO[str],
2087 section: str | None = None,
2088 include_examples: bool = True,
2089 include_descriptions: bool = True,
2090 include_sources: bool = True,
2091 include_env_vars: bool = True,
2092 include_providers: bool = True,
2093 comment_out_everything: bool = False,
2094 hide_sensitive: bool = False,
2095 extra_spacing: bool = True,
2096 only_defaults: bool = False,
2097 show_values: bool = False,
2098 **kwargs: Any,
2099 ) -> None:
2100 """
2101 Write configuration with comments and examples to a file.
2102
2103 :param file: file to write to
2104 :param section: section of the config to write, defaults to all sections
2105 :param include_examples: Include examples in the output
2106 :param include_descriptions: Include descriptions in the output
2107 :param include_sources: Include the source of each config option
2108 :param include_env_vars: Include environment variables corresponding to each config option
2109 :param include_providers: Include providers configuration
2110 :param comment_out_everything: Comment out all values
2111 :param hide_sensitive: Replace sensitive values in the output with "< hidden >"
2112 :param extra_spacing: Add extra spacing before examples and after variables
2113 :param only_defaults: Only include default values when writing the config, not the actual values
2114 """
2115 with self.make_sure_configuration_loaded(with_providers=include_providers):
2116 sources_dict = {}
2117 if include_sources:
2118 sources_dict = self.as_dict(display_source=True)
2119 for section_to_write in self.get_sections_including_defaults():
2120 section_config_description = self.configuration_description.get(section_to_write, {})
2121 if section_to_write != section and section is not None:
2122 continue
2123 if self._has_section_in_any_defaults(section_to_write) or self.has_section(section_to_write):
2124 self._write_section_header(
2125 file, include_descriptions, section_config_description, section_to_write
2126 )
2127 for option in self.get_options_including_defaults(section_to_write):
2128 should_continue, needs_separation = self._write_option_header(
2129 file=file,
2130 option=option,
2131 extra_spacing=extra_spacing,
2132 include_descriptions=include_descriptions,
2133 include_env_vars=include_env_vars,
2134 include_examples=include_examples,
2135 include_sources=include_sources,
2136 section_config_description=section_config_description,
2137 section_to_write=section_to_write,
2138 sources_dict=sources_dict,
2139 )
2140 is_sensitive = self.is_sensitive_option(section_to_write, option)
2141 self._write_value(
2142 file=file,
2143 option=option,
2144 comment_out_everything=comment_out_everything,
2145 needs_separation=needs_separation,
2146 only_defaults=only_defaults,
2147 section_to_write=section_to_write,
2148 hide_sensitive=hide_sensitive,
2149 is_sensitive=is_sensitive,
2150 show_values=show_values,
2151 )
2152 if include_descriptions and not needs_separation:
2153 # extra separation between sections in case last option did not need it
2154 file.write("\n")
2155
2156 @contextmanager
2157 def make_sure_configuration_loaded(self, with_providers: bool) -> Generator[None, None, None]:
2158 """
2159 Make sure configuration is loaded with or without providers.
2160
2161 The context manager will only toggle the `self._use_providers_configuration` flag if `with_providers` is False, and will reset `self._use_providers_configuration` to True after the context block.
2162 Nop for `with_providers=True` as the configuration already loads providers configuration by default.
2163
2164 :param with_providers: whether providers should be loaded
2165 """
2166 if not with_providers:
2167 self._use_providers_configuration = False
2168 # Only invalidate cached properties that depend on _use_providers_configuration.
2169 # Do NOT use invalidate_cache() here — it would also evict expensive provider-discovery
2170 # caches (_provider_metadata_configuration_description, _provider_metadata_config_fallback_default_values)
2171 # that don't depend on this flag.
2172 self._invalidate_provider_flag_caches()
2173 try:
2174 yield
2175 finally:
2176 if not with_providers:
2177 self._use_providers_configuration = True
2178 self._invalidate_provider_flag_caches()