Repository navigation
Conversation
steveahnahn
left a comment
There was a problem hiding this comment.
Verified both paths, _serialize_extended matches canonical BaseSerialization.serialize and round-trips through the real deserializer across nested/empty/list/primitive shapes, and on a live Postgres the pg_temp.encode_extended/decode_extended functions produce the same result and invert cleanly.
nit: the plpgsql path has no test. test_0094_deadline_callback_migration.py builds its engine with _make_engine() (in-memory SQLite), so even the Postgres CI job never executes encode_extended/decode_extended
| 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); |
There was a problem hiding this comment.
I wonder if node = [], whether the entire return value would be None. I guess that's not what you want, right? Should we handle this case here?
There was a problem hiding this 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
| 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); |
There was a problem hiding this comment.
Same fix applied here too — pushed.
Drafted-by: Claude Sonnet 5; reviewed by @seanmuth before posting
| @@ -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) | |||
There was a problem hiding this comment.
Scoping this to nested dicts undersells it. BaseSerialization.deserialize fetches encoded_var[Encoding.VAR] on any non-primitive, non-list value, so a flat kwargs={"k": "v"} and even kwargs={} raise KeyError('__var') exactly like the nested case, and on MySQL/SQLite the pre-fix path wrote no top-level envelope at all so the outer dict raises immediately. Every row 0094 migrated is unreadable, on both backends, not just the ones with nested tags.
Your own diff shows it: test_null_callback_does_not_crash covers a row whose kwargs is {}, and swapping its json.loads for BaseSerialization.deserialize only passes because of this fix. Worth rewording to something like "every deadline callback row migrated by 0094", otherwise an operator whose tags are flat strings reads this release note and concludes they are in the clear.
There was a problem hiding this 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
| _ASYNC_CALLBACK_CLASSNAME = "airflow.sdk.definitions.deadline.AsyncCallback" | ||
|
|
||
|
|
||
| def _serialize_extended(value): |
There was a problem hiding this comment.
0094 is in the released 3.2.0 tag and main is already 3.4.0, so patching it in place only helps deployments upgrading straight from <3.2. Everyone who went through 3.2.x or 3.3.x has already run the broken version, has 100% of their callback rows unreadable, and will never re-run 0094 to pick this up.
Is a repair revision planned as a follow-up? A new migration could reuse these two helpers to re-encode existing rows, and it would need to handle both corrupt shapes: Postgres rows have the outer __var envelope with raw kwargs underneath, MySQL/SQLite rows have no envelope at all. That second shape is why "repair out of band" is not actionable as written -- the jsonb_path_exists query in the description matches on $.**.__var.*, which finds nothing in a MySQL/SQLite row and so reports it clean.
There was a problem hiding this 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
|
@seanmuth are you still working on this one? |
Migration 0094 (replace_deadline_inline_callback_with_fkey) moves inline
deadline callbacks into the callback table, whose data column is ExtendedJSON
(read back via BaseSerialization.deserialize). Both upgrade paths built the
callback data without recursively extended-serializing the callback kwargs:
- Postgres: json_build_object wrapped only the top level, embedding the old
kwargs jsonb raw.
- MySQL/SQLite: wrote the callback dict with no extended wrapping at all.
Any callback whose kwargs contain a nested dict (e.g. metric tags) is then
unreadable -- BaseSerialization.deserialize raises KeyError('__var') and the
scheduler enters CrashLoopBackOff on the deadline-processing loop.
Both upgrade paths now extended-serialize the callback data (a recursive
pg_temp SQL function for the Postgres CTE, a _serialize_extended helper for
the Python path), and the downgrade paths decode kwargs back to raw so an
upgrade/downgrade round-trip is preserved. Deployments already upgraded across
0094 must repair existing callback.data rows out of band.
Closes: apache#69980
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
jsonb_agg over zero rows returns SQL NULL, not '[]'::jsonb, so a callback kwargs value that happens to be an empty list would silently become null on encode (or decode) instead of round-tripping as []. Both pg_temp.encode_extended and decode_extended now COALESCE to '[]'::jsonb, verified against a live Postgres instance. Also widened the newsfragment/PR description's framing: the bug isn't scoped to nested dicts specifically -- any dict-shaped kwargs (flat or even empty) hits the same KeyError, since deserialize requires kwargs itself to be wrapped, not just anything nested inside it. Added a Postgres-backend regression test exercising the real SQL functions directly, closing the gap that let the empty-list bug sit unreviewed for two months -- the existing tests only covered the Python-side helpers used by the MySQL/SQLite path. Co-Authored-By: Claude <noreply@anthropic.com>
Routine lockfile refresh (uv.lock's exclude-newer-span pins to a rolling 4-day window, so this drifts independent of any source change) -- required for this branch's pre-push checks to pass after rebasing onto a newer main. Unrelated to the migration fix itself. Co-Authored-By: Claude <noreply@anthropic.com>
88e738c to
a614565
Compare
|
Yep, still working on it — just let it sit too long. Pushed fixes today for Aaron's and Kaxil's review (a real empty-list encoding bug Aaron caught, plus scope/wording corrections from Kaxil) and rebased onto main. Should be ready for another look. Drafted-by: Claude Sonnet 5; reviewed by @seanmuth before posting |
|
Quickest fix: git fetch upstream main && git rebase upstream/main
rm uv.lock && uv lock
git add uv.lock && git rebase --continue
git push --force-with-leaseAutomated nudge — ignore if you're not ready to rebase. This comment is updated in place on future |
Closes #69980.
What's the problem
Migration
0094_3_2_0_replace_deadline_inline_callback_with_fkeymoves inline deadline callbacks into thecallbacktable, whosedatacolumn isExtendedJSON— read back at runtime viaBaseSerialization.deserialize, which requires every nested dict to be wrapped as{"__type": "dict", "__var": {...}}.Both upgrade paths built
callback.datawithout recursively extended-serializing the callbackkwargs:_upgrade_postgresql):json_build_objectwrapped only the top level, embedding the old callback'skwargsjsonb raw._upgrade_mysql_sqlite): wrote the callback dict via a plainJSONcolumn with no extended wrapping at all.So any callback whose
kwargsis a dict — nested, flat, or even empty ({}) — is stored unreadable, sinceBaseSerialization.deserializerequireskwargsitself to be wrapped, not just anything nested inside it. When the scheduler loads it in the deadline-processing loop,deserializerecurses into the outerdictnode, hits the barekwargsobject, runsencoded_var[Encoding.VAR], and raisesKeyError(<Encoding.VAR: '__var'>)— crashing the scheduler on every loop (CrashLoopBackOff), before it heartbeats, with no OOM and no other logged exception.What's in this PR
BaseSerialization.serialize:pg_temp.encode_extended(jsonb)function used inside the existing batch CTE (keeps the one-round-trip design)._serialize_extended()helper applied to the callback dict.kwargsback to raw (pg_temp.decode_extended/_deserialize_extended), so an upgrade→downgrade round-trip is preserved. Both decoders are lenient, so they also correctly downgrade data written by the pre-fix version of this migration.NULL-corruption edge case in the Postgres encode/decode SQL functions:jsonb_aggover an empty array returns SQLNULL, not'[]'::jsonb, so an empty-list callback value would silently becomenullinstead of round-tripping as[]. Both functions nowCOALESCEto'[]'::jsonb.test_0094_deadline_callback_migration.py): existing assertions now verifycallback.dataround-trips through the realBaseSerialization; added a nested-kwargs/tagsregression test, an encode/decode helper round-trip test, and a Postgres-backend test (@pytest.mark.backend("postgres")) covering the empty-list fix directly against the real SQL functions.Already-affected deployments
This hotfix protects deployments that have not yet run
0094. Deployments already upgraded across0094with nested-dict callback kwargs have latent corruptcallback.datarows (a scheduler crash that fires when such a deadline becomes due); those rows must be repaired out of band — e.g. re-encode the affected rows:🤖 Opened with Claude Code (model: Claude Opus 4.8,
claude-opus-4-8[1m]).