Repository navigation
Pass AssetStateStoreAccessors into response_check_path Callable in HttpEventTrigger - #72287
sundeep8967 wants to merge 2 commits into
Conversation
c65e94e to
fb07662
Compare
potiuk
left a comment
There was a problem hiding this comment.
Thanks for picking this up — wiring asset_state_store into the response check is the right direction, and replacing the Variable.set() workaround the docs currently recommend is a real improvement. But the way the second argument is detected changes behaviour for callables that work today, and the attribute it reads does not exist on most Airflow versions this provider supports.
Major: arity-based dispatch breaks existing multi-parameter callables
Description doesn't match code: PR description describes something different from what the code actually does.
sig = inspect.signature(response_check)
params = list(sig.parameters.values())
if len(params) >= 2 or any(
p.kind in (inspect.Parameter.VAR_POSITIONAL, inspect.Parameter.VAR_KEYWORD) for p in params
):
check = await response_check(response, self.asset_state_store)The description says the store is passed "if the callable accepts it", but the code passes it positionally whenever the signature has any second parameter. I ran the condition against a few shapes:
| Callable | store passed? | result |
|---|---|---|
async def check(response) |
no | unchanged |
async def check(response, threshold=5) |
yes | threshold silently receives an AssetStateStoreAccessors (or None) |
async def check(response, *, strict=True) |
yes | TypeError: takes 1 positional argument but 2 were given |
async def check(response, **kwargs) |
yes | store lands in kwargs |
Both middle rows are callables that work on main today. Suggestion: match by name and pass it as a keyword — that is what "accepts it" means, and it cannot collide with a user's own parameters:
sig = inspect.signature(response_check)
accepts_store = "asset_state_store" in sig.parameters or any(
p.kind is inspect.Parameter.VAR_KEYWORD for p in sig.parameters.values()
)
if accepts_store:
check = await response_check(response, asset_state_store=self.asset_state_store)
else:
check = await response_check(response)(*args-only callables then keep the single-argument call, which is the safe default.)
Major: self.asset_state_store only exists on Airflow >= 3.3.0
Use
version_compat.pypatterns for cross-version compatibility. —providers/AGENTS.md
asset_state_store is set in BaseEventTrigger.__init__ (airflow-core/src/airflow/triggers/base.py:322), which first shipped in Airflow 3.3.0. This provider declares apache-airflow>=2.11.0 (providers/http/pyproject.toml:62), so on Airflow 3.0–3.2 a two-parameter callable raises AttributeError, which run() swallows in its except Exception, logs as "Poll failed", and retries max_consecutive_failures times before giving up — a confusing failure on a supported version.
The providers that already consume the store guard for exactly this:
store = getattr(self, "asset_state_store", None)(providers/amazon/.../triggers/kinesis.py:153, providers/apache/iceberg/.../triggers/iceberg.py:125). Please do the same here, or gate on AIRFLOW_V_3_3_PLUS from airflow.providers.common.compat.version_compat, which this module already imports from. Note the value is also None on 3.3+ when the trigger is not inflated for an asset watcher, so the docstring should tell users to expect None.
Smaller observations
- Docs out of sync (
providers/AGENTS.md: "Keepprovider.yamlmetadata, docs, and tests in sync."): neither the:param response_check_path:docstring norproviders/http/docs/triggers.rstmentions the new argument, and the docs page still recommends theVariable.set()workaround in "Important Notes" that this PR is meant to replace. Please document the optionalasset_state_storeparameter (and itsNonecases) and update the example. - Test coverage of the new condition: the
*args/**kwargsbranch is untested, and so is the case that matters most for the first point — a callable with a second defaulted or keyword-only parameter that must keep receiving one argument. PerAGENTS.md("Use@pytest.mark.parametrizefor multiple similar inputs"), those plus the two existing tests fit naturally in one parametrized test over(callable, expected_call_args). - Sentinel type in the test:
event_trigger.asset_state_store = "mock_store"— prefermock.create_autospec(AssetStateStoreAccessors, instance=True)so the test documents what the callable really receives; a bare string would not catch the code calling a method on the wrong object. - The PR body dropped the template, including the "Was generative AI tooling used to co-author this PR?" section — please restore it and answer it either way.
- The main CI workflow has not run on this head yet (only the
WIPcheck is present), so nothing here is verified by CI; a maintainer will approve the workflow run once the points above are addressed.
This review was drafted by an AI-assisted tool and
confirmed by an Apache Airflow maintainer. After you've
addressed the points above and pushed an update, an Apache Airflow
maintainer — a real person — will take the next look
at the PR. The findings cite the project's review criteria;
if you think one of them is mis-applied, please reply on the
PR and a maintainer will weigh in.More on how Apache Airflow handles maintainer review:
contributing-docs/05_pull_requests.rst.
0960259 to
5bf04b7
Compare
|
Thanks for the thorough review @potiuk! I have updated the implementation and docs to address all points:
|
5bf04b7 to
1d880d0
Compare
There was a problem hiding this comment.
Both points from the last round are properly fixed — the keyword-based dispatch even handles the positional-only case I hadn't called out, and the getattr guard matches the kinesis/iceberg pattern exactly. Two new issues came in with the revision, though: the test import breaks the pre-3.3 compatibility runs, and the docs example teaches users to block the triggerer's event loop.
Major: the new test import breaks the whole module on Airflow < 3.3 (providers/http/tests/unit/http/triggers/test_http.py:39)
from airflow.sdk.execution_time.context import AssetStateStoreAccessorsAssetStateStoreAccessors first shipped in Airflow 3.3.0. PROVIDERS_COMPATIBILITY_TESTS_MATRIX (dev/breeze/src/airflow_breeze/global_constants.py:872) runs http provider unit tests against 2.11.1, 3.0.6, 3.1.8, 3.2.2 and 3.3.2, and http is not in any remove-providers list — so on four of those five this import raises at collection time and takes down every test in the file, including the ~20 that have nothing to do with this change.
The Iceberg trigger hit the same wall and guards it by importing inside the version-gated test (providers/apache/iceberg/tests/unit/apache/iceberg/triggers/test_iceberg.py:232):
@pytest.mark.skipif(not AIRFLOW_V_3_3_PLUS, reason="asset_state_store arrived in Airflow 3.3.0")
async def test_reads_the_watermark_through_the_real_asset_state_store():
from airflow.sdk.execution_time.context import AssetStateStoreAccessorsAIRFLOW_V_3_3_PLUS comes from tests_common.test_utils.version_compat. Note the rest of that file's tests use a plain MagicMock() store, which needs no import at all — that works for every shape in test_run_response_check_callable_shapes except the one asserting identity against the autospec.
Major: the docs example blocks the triggerer's event loop (providers/http/docs/triggers.rst:73)
aget— "Async version ofgetthat awaits instead of blocking the event loop." —AssetStateStoreAccessors,task-sdk/src/airflow/sdk/execution_time/context.py:906
async def check_github_api_response(response, asset_state_store=None):
...
previous_release_id = asset_state_store.get("release_id", None) # <- blocking
...
asset_state_store.set("release_id", release_id) # <- blockingAssetStateStoreAccessor.get/.set do a synchronous SUPERVISOR_COMMS.send(...) round-trip (context.py:759). The response_check_path callable runs on the triggerer's shared event loop, so this example — the thing users copy — stalls every other trigger in the process for the duration of each round-trip.
One wrinkle worth documenting rather than glossing: aget/aset arrived in #72127 and are not in any released version yet (3.3.0 and 3.3.1 expose only the blocking API). The Kinesis trigger handles exactly this (providers/amazon/src/airflow/providers/amazon/aws/triggers/kinesis.py:155-159):
# aget/aset landed in Airflow 3.3.2; 3.3.0 and 3.3.1 only expose the blocking API.
if hasattr(store, "aget"):
checkpoint = await store.aget(key, default={})
else:
checkpoint = await asyncio.to_thread(store.get, key, default={})Mirroring that in the example (and saying why) gives users something safe to copy on every version where the feature exists.
Smaller observations
See inline comments on providers/http/docs/triggers.rst:100 and providers/http/tests/unit/http/triggers/test_http.py:394.
Worth a second look from
This change touches the http trigger and the asset-state-store contract; folks with the most context here:
@amoghrajesh— 12 of the last 30 commits ontask-sdk/src/airflow/sdk/execution_time/context.py, where the accessor lives, plus a commit on this trigger file (committer)@karenbraganz— authoredHttpEventTriggerand theVariable.set()docs example this PR replaces@1fanwang— wrote the Iceberg trigger's asset-state-store integration, the closest precedent for both findings above
None of them have been notified — asking any of them for an extra pass is the maintainer's call, and optional.
One process note: the main CI workflow still hasn't run on this head (only the WIP and Mergeable bot checks are present), so none of this is verified by CI yet.
This review was drafted by an AI-assisted tool and
confirmed by an Apache Airflow maintainer. After you've
addressed the points above and pushed an update, an Apache Airflow
maintainer — a real person — will take the next look
at the PR. The findings cite the project's review criteria;
if you think one of them is mis-applied, please reply on the
PR and a maintainer will weigh in.More on how Apache Airflow handles maintainer review:
contributing-docs/05_pull_requests.rst.
Drafted-by: Claude Code (Opus 5); reviewed by @potiuk before posting
…tpEventTrigger HttpEventTrigger inherits from BaseEventTrigger which provides asset_state_store to read and persist state for an Asset. This update inspects the response_check callable signature and passes asset_state_store as a keyword argument if the callable accepts it (or **kwargs), avoiding positional argument collisions. It also safely guards against AttributeError on Airflow versions earlier than 3.3.0 via getattr(self, "asset_state_store", None), updates docs and examples with non-blocking aget/aset calls and downstream task consumption, and expands tests with mock AssetStateStoreAccessors. Fixes apache#70074 Signed-off-by: sundeep8967 <sundeep8967@gmail.com>
1b8694b to
a940ecf
Compare
|
Thanks @potiuk for the review and catching these points! I have updated the PR to resolve both:
|
potiuk
left a comment
There was a problem hiding this comment.
Thanks — both points from last round are properly addressed: the AssetStateStoreAccessors import is now gated inside the 3.3+ test, and the docs example uses aget/aset with the asyncio.to_thread fallback, exactly the Kinesis shape.
Since the main CI workflow still hasn't run on this head, I applied the PR locally and ran the new tests against it. Result: 6 failed, 2 passed.
Blocking — every case of test_run_response_check_callable_shapes fails
The parametrized checks are plain lambdas, and _run_response_check rejects any non-coroutine callable before it looks at the signature:
providers/http/src/airflow/providers/http/triggers/http.py:398: in _run_response_check
E AirflowException: The response_check callable is not asynchronous.
They need to be async def functions — module-level helpers parametrized by reference work well, and that also removes the need for eval(...) in the positional-only case: async def _positional_only(resp, asset_state_store="default_val", /) is valid syntax on every Python version Airflow supports. Details inline.
Major — the docs example's downstream task never sees the store
context["asset_state_store"] is only populated when the task has concrete inlets or outlets (task_runner.py:350); schedule=asset on the Dag does not count. As written, print_airflow_release_info always gets None and prints Unknown has been released. Declaring the asset as an inlet fixes it — @task(inlets=[asset]) — and then store.get(...) resolves to the single accessor. (This follows up on my own suggestion from last round, which was missing that detail.)
Minor
On Airflow < 3.3 asset_state_store is None, so the example's check skips deduplication and returns True on every poll — the Dag would trigger every 60 seconds. The old Variable-based example deduplicated on every version the provider supports. Either keep a Variable fallback in that branch, or say plainly in the note that the example requires Airflow 3.3+.
This review was drafted by an AI-assisted tool and
confirmed by an Apache Airflow maintainer. After you've
addressed the points above and pushed an update, an Apache Airflow
maintainer — a real person — will take the next look
at the PR. The findings cite the project's review criteria;
if you think one of them is mis-applied, please reply on the
PR and a maintainer will weigh in.More on how Apache Airflow handles maintainer review:
contributing-docs/05_pull_requests.rst.
Drafted-by: Claude Code (Opus 5); reviewed by @potiuk before posting
| "check", | ||
| [ | ||
| pytest.param( | ||
| lambda resp: resp == "ok", |
There was a problem hiding this comment.
Blocking: these are sync lambdas, so _run_response_check raises before it gets to the signature check, and all six cases fail:
providers/http/src/airflow/providers/http/triggers/http.py:398: in _run_response_check
E AirflowException: The response_check callable is not asynchronous.
Module-level async def helpers, parametrized by reference, fix it. They also replace the eval(...) below, since positional-only parameters are ordinary syntax:
async def _positional_only(resp, asset_state_store="default_val", /):
return resp == "ok" and asset_state_store == "default_val"| @@ -108,10 +113,10 @@ Here's an example of using the ``HttpEventTrigger`` in an ``AssetWatcher`` to mo | |||
| @dag(start_date=datetime.datetime(2024, 10, 1), schedule=asset, catchup=False) | |||
| def check_airflow_releases(): | |||
| @task() | |||
There was a problem hiding this comment.
Major: context["asset_state_store"] is only populated for tasks with concrete inlets or outlets (task_runner.py:350), and schedule=asset doesn't count. As written, store is always None here and the task prints Unknown has been released. Declaring the asset as an inlet fixes it:
| @task() | |
| @task(inlets=[asset]) |
Closes: #70074
Description
With the event-driven asset enhancements in
BaseEventTrigger,asset_state_storeis available on Airflow >= 3.3.0 to read and persist state for an asset across evaluations.This PR:
_run_response_checkinHttpEventTrigger(providers/http/src/airflow/providers/http/triggers/http.py) to inspect theresponse_checkcallable signature and passasset_state_storeby keyword argument if the callable accepts it (or**kwargs), preventing positional argument collisions with default or keyword-only arguments.store = getattr(self, "asset_state_store", None)to preventAttributeErroron Airflow versions earlier than 3.3.0 or when running outside an asset watcher context.providers/http/docs/triggers.rstand trigger docstrings to document the optionalasset_state_storeparameter (and when it can beNone), updating the release watcher example and replacing the outdatedVariable.set()note.mock.create_autospec(AssetStateStoreAccessors, instance=True), as well as a test for execution whenasset_state_storeis absent.Was generative AI tooling used to co-author this PR?
{pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.