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]