seanmuth commented on code in PR #69985:
URL: https://github.com/apache/airflow/pull/69985#discussion_r4134019060


##########
airflow-core/src/airflow/migrations/versions/0094_3_2_0_replace_deadline_inline_callback_with_fkey.py:
##########
@@ -54,12 +54,102 @@
 _ASYNC_CALLBACK_CLASSNAME = "airflow.sdk.definitions.deadline.AsyncCallback"
 
 
+def _serialize_extended(value):
+    """
+    Encode ``value`` in Airflow's extended-JSON format, matching 
``BaseSerialization.serialize``.
+
+    ``callback.data`` is an ``ExtendedJSON`` column, so on read the runtime 
runs
+    ``BaseSerialization.deserialize``, which requires every nested dict to be 
wrapped as
+    ``{"__type": "dict", "__var": {...}}``. We must write that same wrapping 
here (recursively)
+    rather than embedding raw nested dicts -- otherwise deserialize raises 
``KeyError('__var')``
+    and crashes the scheduler. The dict/list/primitive subset below is all 
that callback data
+    contains; the logic is inlined so the migration stays independent of 
runtime serialization code.
+    """
+    if isinstance(value, dict):
+        return {"__type": "dict", "__var": {str(k): _serialize_extended(v) for 
k, v in value.items()}}
+    if isinstance(value, list):
+        return [_serialize_extended(v) for v in value]
+    return value
+
+
+# Recursive extended-JSON encoder as a session-local (auto-dropped) SQL 
function, so the
+# Postgres CTE path can wrap nested callback kwargs the same way 
``_serialize_extended`` does.
+_PG_ENCODE_EXTENDED_DDL = dedent("""
+    CREATE OR REPLACE FUNCTION pg_temp.encode_extended(node jsonb) RETURNS 
jsonb AS $$
+    DECLARE k text; v jsonb; out jsonb := '{}'::jsonb;
+    BEGIN
+      IF jsonb_typeof(node) = 'object' THEN
+        FOR k, v IN SELECT * FROM jsonb_each(node) LOOP
+          out := out || jsonb_build_object(k, pg_temp.encode_extended(v));
+        END LOOP;
+        RETURN jsonb_build_object('__type', 'dict', '__var', out);
+      ELSIF jsonb_typeof(node) = 'array' THEN
+        RETURN (SELECT jsonb_agg(pg_temp.encode_extended(e)) FROM 
jsonb_array_elements(node) e);

Review Comment:
   Good catch — confirmed this is a real bug: `jsonb_agg` over zero rows 
returns SQL `NULL`, not `'[]'::jsonb`, so `node = []` would silently become 
`null` instead of round-tripping as `[]`. Fixed with `COALESCE(..., 
'[]'::jsonb)`, verified against a live Postgres instance, and added a 
regression test (`TestMigration0094PostgresEncodeDecodeHelpers`) that exercises 
the real SQL functions directly — pushed.
   
   ---
   Drafted-by: Claude Sonnet 5; reviewed by @seanmuth before posting



##########
airflow-core/src/airflow/migrations/versions/0094_3_2_0_replace_deadline_inline_callback_with_fkey.py:
##########
@@ -54,12 +54,102 @@
 _ASYNC_CALLBACK_CLASSNAME = "airflow.sdk.definitions.deadline.AsyncCallback"
 
 
+def _serialize_extended(value):
+    """
+    Encode ``value`` in Airflow's extended-JSON format, matching 
``BaseSerialization.serialize``.
+
+    ``callback.data`` is an ``ExtendedJSON`` column, so on read the runtime 
runs
+    ``BaseSerialization.deserialize``, which requires every nested dict to be 
wrapped as
+    ``{"__type": "dict", "__var": {...}}``. We must write that same wrapping 
here (recursively)
+    rather than embedding raw nested dicts -- otherwise deserialize raises 
``KeyError('__var')``
+    and crashes the scheduler. The dict/list/primitive subset below is all 
that callback data
+    contains; the logic is inlined so the migration stays independent of 
runtime serialization code.
+    """
+    if isinstance(value, dict):
+        return {"__type": "dict", "__var": {str(k): _serialize_extended(v) for 
k, v in value.items()}}
+    if isinstance(value, list):
+        return [_serialize_extended(v) for v in value]
+    return value
+
+
+# Recursive extended-JSON encoder as a session-local (auto-dropped) SQL 
function, so the
+# Postgres CTE path can wrap nested callback kwargs the same way 
``_serialize_extended`` does.
+_PG_ENCODE_EXTENDED_DDL = dedent("""
+    CREATE OR REPLACE FUNCTION pg_temp.encode_extended(node jsonb) RETURNS 
jsonb AS $$
+    DECLARE k text; v jsonb; out jsonb := '{}'::jsonb;
+    BEGIN
+      IF jsonb_typeof(node) = 'object' THEN
+        FOR k, v IN SELECT * FROM jsonb_each(node) LOOP
+          out := out || jsonb_build_object(k, pg_temp.encode_extended(v));
+        END LOOP;
+        RETURN jsonb_build_object('__type', 'dict', '__var', out);
+      ELSIF jsonb_typeof(node) = 'array' THEN
+        RETURN (SELECT jsonb_agg(pg_temp.encode_extended(e)) FROM 
jsonb_array_elements(node) e);
+      END IF;
+      RETURN node;
+    END;
+    $$ LANGUAGE plpgsql;
+""")
+
+
+def _deserialize_extended(value):
+    """
+    Inverse of :func:`_serialize_extended`: unwrap extended-JSON back to plain 
values.
+
+    Used by ``downgrade`` to rebuild the old inline callback ``kwargs`` (which 
were stored
+    raw). Lenient: already-raw dicts (e.g. produced by the pre-fix version of 
this migration)
+    pass through unchanged, so downgrade is correct regardless of which 
version upgraded.
+    """
+    if isinstance(value, dict):
+        if "__type" in value and "__var" in value:
+            if value["__type"] == "dict" and isinstance(value["__var"], dict):
+                return {k: _deserialize_extended(v) for k, v in 
value["__var"].items()}
+            return value
+        return {k: _deserialize_extended(v) for k, v in value.items()}
+    if isinstance(value, list):
+        return [_deserialize_extended(v) for v in value]
+    return value
+
+
+# SQL inverse of ``pg_temp.encode_extended`` for the Postgres downgrade path. 
Lenient in the
+# same way as ``_deserialize_extended`` so it handles data written by either 
version of upgrade.
+_PG_DECODE_EXTENDED_DDL = dedent("""
+    CREATE OR REPLACE FUNCTION pg_temp.decode_extended(node jsonb) RETURNS 
jsonb AS $$
+    DECLARE k text; v jsonb; out jsonb := '{}'::jsonb;
+    BEGIN
+      IF jsonb_typeof(node) = 'object' THEN
+        IF node ? '__type' AND node ? '__var' THEN
+          IF node->>'__type' = 'dict' AND jsonb_typeof(node->'__var') = 
'object' THEN
+            FOR k, v IN SELECT * FROM jsonb_each(node->'__var') LOOP
+              out := out || jsonb_build_object(k, pg_temp.decode_extended(v));
+            END LOOP;
+            RETURN out;
+          END IF;
+          RETURN node;
+        ELSE
+          FOR k, v IN SELECT * FROM jsonb_each(node) LOOP
+            out := out || jsonb_build_object(k, pg_temp.decode_extended(v));
+          END LOOP;
+          RETURN out;
+        END IF;
+      ELSIF jsonb_typeof(node) = 'array' THEN
+        RETURN (SELECT jsonb_agg(pg_temp.decode_extended(e)) FROM 
jsonb_array_elements(node) e);

Review Comment:
   Same fix applied here too — pushed.
   
   ---
   Drafted-by: Claude Sonnet 5; reviewed by @seanmuth before posting



##########
airflow-core/newsfragments/69985.bugfix.rst:
##########
@@ -0,0 +1 @@
+Fix scheduler crash from deadline callbacks whose kwargs contain a nested 
dict. Migration ``0094`` (``replace_deadline_inline_callback_with_fkey``) built 
``callback.data`` with nested kwargs left unencoded, so 
``BaseSerialization.deserialize`` raised ``KeyError`` on read and the scheduler 
entered CrashLoopBackOff. Both the Postgres and MySQL/SQLite upgrade paths now 
extended-serialize nested kwargs (and the downgrade path decodes them back). 
Deployments already upgraded across ``0094`` must repair existing 
``callback.data`` rows out of band. (#69980)

Review Comment:
   You're right, this undersold it — reworded the newsfragment and PR 
description: any dict-shaped `kwargs` (flat, or even empty `{}`) hits the same 
`KeyError`, not just ones with further nested structure, since `deserialize` 
requires `kwargs` itself to be wrapped regardless of what's inside it.
   
   ---
   Drafted-by: Claude Sonnet 5; reviewed by @seanmuth before posting



##########
airflow-core/src/airflow/migrations/versions/0094_3_2_0_replace_deadline_inline_callback_with_fkey.py:
##########
@@ -54,12 +54,102 @@
 _ASYNC_CALLBACK_CLASSNAME = "airflow.sdk.definitions.deadline.AsyncCallback"
 
 
+def _serialize_extended(value):

Review Comment:
   Agreed this only helps deployments that haven't run `0094` yet — that's what 
the PR body's "Already-affected deployments" section covers: already-upgraded 
deployments need the out-of-band repair (SQL provided there), since patching an 
already-applied migration can't retroactively fix rows it already wrote. This 
split (hotfix-in-place + out-of-band repair, no new repair migration) was what 
Daniel Standish suggested on the original issue. Let me know if you think that 
approach needs revisiting.
   
   ---
   Drafted-by: Claude Sonnet 5; reviewed by @seanmuth before posting



-- 
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