beaglebay opened a new issue, #70398:
URL: https://github.com/apache/airflow/issues/70398

   ## Apache Airflow version
   
   3.0.6
   
   ## If "Other Airflow 2 version" selected, which one?
   
   _No response_
   
   ## What happened?
   
   Deferred `DbtCloudRunJobOperator` tasks fail almost immediately with a 
misleading `DbtCloudJobRunException: Job run <id> has failed.` even though the 
real dbt Cloud job is still running (and later completes successfully, 
sometimes hours later). The task fails within ~10 seconds of being deferred — 
far too fast to reflect a real terminal job state for a job that in fact took 
several more hours to finish.
   
   I traced this down to a real bug in Airflow core (Task SDK), not in the dbt 
Cloud provider itself, and not specific to any one secrets backend 
implementation. Full trace below.
   
   ### Environment
   
   - Deployment: Amazon MWAA
   - `secrets.backend`: 
`airflow.providers.amazon.aws.secrets.secrets_manager.SecretsManagerBackend`
   - `secrets.backend_kwargs`: 
`{"config_prefix":"airflow/config","connections_prefix":"airflow/connections","variables_prefix":"airflow/variables"}`
   - Connection `dbt_cloud` (conn_type `dbt_cloud`) has `password` populated 
(used by the provider as the dbt Cloud API token) and no `extra`.
   - The task is `DbtCloudRunJobOperator(..., deferrable=True)` using the 
connection above.
   
   ### Observed log sequence
   
   **Worker (synchronous `execute()`, before deferring) — succeeds fine:**
   ```
   {"logger": "airflow.hooks.base", "event": "Connection Retrieved 
'dbt_cloud'", "level": "info"}
   {"logger": "...DbtCloudHook", "event": "Getting the status of job run 
<job_run_id>.", "level": "info"}
   {"logger": "...DbtCloudHook", "event": "Current status of job run 
<job_run_id>: QUEUED", "level": "info"}
   {"event": "Pausing task as DEFERRED.", "level": "info"}
   ```
   
   **Triggerer (async `DbtCloudRunJobTrigger.run()` poll cycle) — crashes:**
   ```
   {"event": "trigger <dag_id>/<run_id>/dbt_cloud_task/-1/1 (ID <n>) starting", 
"level": "info"}
   {"logger": "...DbtCloudHook", "event": "Getting the status of job run 
<job_run_id>.", "level": "info"}
   {
     "logger": "task",
     "error_detail": [{
       "exc_type": "RuntimeError",
       "exc_value": "You cannot use AsyncToSync in the same thread as an async 
event loop - just await the async function directly.",
       "frames": [
         {"filename": ".../airflow/sdk/execution_time/context.py", "lineno": 
131, "name": "_get_connection"},
         {"filename": ".../airflow/secrets/base_secrets.py", "lineno": 79, 
"name": "get_connection"},
         {"filename": ".../airflow/secrets/base_secrets.py", "lineno": 65, 
"name": "deserialize_connection"},
         {"filename": ".../airflow/models/connection.py", "lineno": 557, 
"name": "from_json"},
         {"filename": "<string>", "lineno": 4, "name": "__init__"},
         {"filename": ".../sqlalchemy/orm/state.py", "lineno": 481, "name": 
"_initialize_instance"},
         {"filename": ".../sqlalchemy/util/langhelpers.py", "lineno": 70, 
"name": "__exit__"},
         {"filename": ".../sqlalchemy/util/compat.py", "lineno": 211, "name": 
"raise_"},
         {"filename": ".../sqlalchemy/orm/state.py", "lineno": 479, "name": 
"_initialize_instance"},
         {"filename": ".../airflow/models/connection.py", "lineno": 177, 
"name": "__init__"},
         {"filename": ".../airflow/sdk/execution_time/secrets_masker.py", 
"lineno": 134, "name": "mask_secret"},
         {"filename": ".../airflow/jobs/triggerer_job_runner.py", "lineno": 
740, "name": "send"},
         {"filename": ".../asgiref/sync.py", "lineno": 186, "name": "__call__"}
       ]
     }],
     "event": "Task failed with exception",
     "level": "error"
   }
   {"event": "Unable to retrieve connection from secrets backend 
(SecretsManagerBackend). Checking subsequent secrets backend.", "level": 
"error"}
   {"result": "TriggerEvent<{'status': 'error', 'message': 'You cannot use 
AsyncToSync in the same thread as an async event loop - just await the async 
function directly.', 'run_id': <job_run_id>}>", "event": "Trigger fired event", 
"level": "info"}
   {"event": "trigger completed", "level": "info"}
   ```
   
   This `TriggerEvent(status='error', ...)` is then consumed by 
`DbtCloudRunJobOperator.execute_complete()`, which raises 
`DbtCloudJobRunException: Job run <job_run_id> has failed.` — a message that 
implies the *dbt Cloud job* failed, when in reality Airflow never actually 
observed its terminal state at all.
   
   **Confirmed independently via the dbt Cloud Admin API** that the underlying 
run continued after Airflow marked the task failed, and later reached a real 
terminal state (`Success`) several hours afterward — proving the "has failed" 
exception was not based on any actual observed job outcome.
   
   ### Root cause (traced against `apache/airflow` tag `3.0.6`)
   
   This is not a dbt Cloud provider bug, and not specific to 
`SecretsManagerBackend`. It's in `_get_connection()`, 
`task-sdk/src/airflow/sdk/execution_time/context.py`:
   
   ```python
   def _get_connection(conn_id: str) -> Connection:
       backends = ensure_secrets_backend_loaded()
       for secrets_backend in backends:
           try:
               conn = secrets_backend.get_connection(conn_id=conn_id)
               if conn:
                   if conn.password:
                       mask_secret(conn.password)   # <-- crashes here
                   if conn.extra:
                       mask_secret(conn.extra)
                   return conn
           except Exception:
               log.exception(
                   "Unable to retrieve connection from secrets backend (%s). "
                   "Checking subsequent secrets backend.",
                   type(secrets_backend).__name__,
               )
       ...
       msg = SUPERVISOR_COMMS.send(GetConnection(conn_id=conn_id))  # <-- same 
crash, uncaught here
   ```
   
   1. `secrets_backend.get_connection(conn_id)` **succeeds** — the connection 
is fetched and deserialized fine.
   2. `mask_secret(conn.password)` is called next. `mask_secret()` 
(`task-sdk/src/airflow/sdk/execution_time/secrets_masker.py:132-134`) does:
      ```python
      if comms := getattr(task_runner, "SUPERVISOR_COMMS", None):
          comms.send(MaskSecret(value=secret, name=name))
      ```
   3. Inside the triggerer, `SUPERVISOR_COMMS` is the triggerer's comms decoder 
(per `airflow/jobs/triggerer_job_runner.py`). Calling `.send()` on it from 
inside a coroutine that is *already* running the active event loop conflicts 
with the sync→async bridge (`asgiref.sync.AsyncToSync`) it uses internally, 
raising the `RuntimeError`.
   4. This exception is caught by the `except Exception` in the loop (logged as 
"Checking subsequent secrets backend"), but since only one backend is 
configured, the loop ends. Execution falls through to the **unguarded** 
fallback call `SUPERVISOR_COMMS.send(GetConnection(conn_id=conn_id))` at line 
160 — the exact same `.send()` primitive, called from the exact same broken 
async context — which raises the identical `RuntimeError`, this time uncaught, 
propagating up into the trigger's `run()` coroutine and surfacing as a generic 
`TriggerEvent(status='error')`.
   
   Key point: **this has nothing to do with which secrets backend is 
configured, or which Connection field holds the secret.** Any deferred operator 
that (a) uses a Connection with `password` or `extra` populated, and (b) 
re-fetches that connection from inside the triggerer's async `run()` loop (as 
`DbtCloudRunJobTrigger` does on every poll, via `provide_account_id`'s 
`sync_to_async(self.get_connection)(...)` call), will hit this. It happens to 
surface here via `apache-airflow-providers-dbt-cloud`, but the defect is in 
Airflow core / Task SDK, not the provider.
   
   ### Related reports
   
   - #53447 — same `RuntimeError: You cannot use AsyncToSync...` signature, 
reported for `BigQueryHook` in a custom trigger (Airflow 3.0.3). Closed as 
"custom trigger code," but a later comment on that thread (2026-02-19) points 
out the identical crash occurs inside the **official** 
`DbtCloudRunJobTrigger.run()` — this wasn't spun out into its own issue at the 
time.
   - #64213 — same signature via `TriggerCommsDecoder.send()` cross-loop 
conflict, this time through the amazon provider's async AWS hook. Closed by the 
reporter citing possible triggerer instability rather than a confirmed 
deterministic defect, due to lack of an isolated repro.
   - #69739 (open) — adds trigger-event payload validation to 
`DbtCloudRunJobOperator`/`DbtCloudJobRunSensor`. Useful, but addresses a 
different symptom (malformed/unrecognized event payloads) and does not touch 
this root cause.
   
   This report differs from the above in that the trace above shows the exact 
same failure mode reproduced independently three times in the same ~20 second 
window (three separate deferred `DbtCloudRunJobOperator` tasks polling around 
the same time), each with a matching stack trace ending in 
`secrets_masker.py:134 mask_secret` → `triggerer_job_runner.py:740 send` → 
`asgiref/sync.py:186 AsyncToSync.__call__`, and independently confirmed against 
the downstream API that the reported "job run has failed" was not based on a 
real observed terminal state for at least two of the three runs.
   
   ## What you think should happen instead?
   
   At minimum, `mask_secret()` failures inside `_get_connection()` should not 
be allowed to abort connection retrieval / masquerade as the underlying 
job/task failing — masking is a logging side-effect and shouldn't be able to 
fail the primary operation. More fundamentally, `SUPERVISOR_COMMS.send()` 
(specifically the triggerer's comms decoder implementation) needs to correctly 
bridge calls made from within its own already-running event loop (e.g. via 
`run_coroutine_threadsafe` against the stored main loop) instead of raising 
when `AsyncToSync` detects it's already inside that loop.
   
   ## How to reproduce
   
   1. Configure any secrets backend (e.g. `SecretsManagerBackend`) with a 
Connection that has `password` and/or `extra` populated.
   2. Use that connection with a deferrable operator whose trigger re-fetches 
the connection during its `run()` poll loop (e.g. 
`DbtCloudRunJobOperator(deferrable=True)`, whose trigger calls 
`provide_account_id`'s `sync_to_async(self.get_connection)(...)` on each status 
check).
   3. Trigger the DAG and let the task defer.
   4. Observe the triggerer log for the task: it fails within seconds with 
`RuntimeError: You cannot use AsyncToSync in the same thread as an async event 
loop`, surfaced by the operator as a misleading `DbtCloudJobRunException: Job 
run <id> has failed.`
   5. Separately query the actual job status via the downstream API — it will 
still be in progress or will complete successfully well after Airflow reported 
failure.
   
   ## Operating System
   
   Amazon Linux 2 (MWAA managed runtime, Python 3.12)
   
   ## Versions of Apache Airflow Providers
   
   ```
   apache-airflow-providers-amazon==9.12.0
   apache-airflow-providers-dbt-cloud==4.4.2
   asgiref==3.9.1
   ```
   (per the `constraints-3.0.6/constraints-3.12.txt` constraint file; MWAA 
pre-installs/version-locks these and does not allow pinning them independently.)
   
   ## Deployment
   
   Amazon Managed Workflows for Apache Airflow (MWAA)
   
   ## Deployment details
   
   MWAA environment running Airflow 3.0.6 on Python 3.12, with 
`secrets.backend` configured to `SecretsManagerBackend`. Same crash signature 
reproduced across three separate DAGs (each a single deferrable 
`DbtCloudRunJobOperator` task) within the same ~20 second window.
   
   ## Anything else?
   
   Happy to provide additional sanitized log excerpts if useful. I did not 
include specific connection IDs, hostnames, account identifiers, or credential 
values in this report.
   
   ## Are you willing to submit PR?
   
   - [ ] Yes I am willing to submit a PR!
   
   ## Code of Conduct
   
   - [x] I agree to follow this project's [Code of 
Conduct](https://github.com/apache/airflow/blob/main/CODE_OF_CONDUCT.md)
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to