This is an automated email from the ASF dual-hosted git repository.
jason810496 pushed a change to branch
jason/lang-sdk-e2e/03c-coordinator-dag-parsing
in repository https://gitbox.apache.org/repos/asf/airflow.git
omit d21aea99762 Say dag_bundle_name is only needed for two coordinators of
one class
omit 39fcab0bc98 Drop serves_bundle so get_parsed_bundles is the only answer
omit 27a987d2641 Check Dag file claims per bundle on the parse path
omit f27fb2d3de8 Import the coordinator manager at the top of the importer
base
omit 43d117ddbca Let a coordinator parse the Dag files of the bundles it
serves
omit 6a995f27589 Sync the Java SDK copy of the Dag schema
omit ad6eae0144f Serialize the SDK Dag in tests that write a dag_maker Dag
omit f1f5d581143 Note in the parsing ADR that runtimes may omit
config-backed fields
omit d3effc9c25d Check SDK Dags with the Dag processor's validation in
conformance
omit 2482920cb7e Reject malformed task entries in a serialized Dag
omit d7e9d5f5e72 Fill a serialized Dag's unset settings from the Airflow
config
omit 0834b7149ac Check that a serialized Dag can be stored and loaded
omit 7c112fe041f Move the Dag file processor's shared plumbing into a base
class
add a3e07159a0b Add Dag ID filtering to Assets search (#70971)
add d3b9a0dd2b3 Add KafkaSharedStreamProducer and KafkaSharedStreamTrigger
(#68625)
add ae6139d58a8 Fix DecreasingPriorityStrategy example by initializing
try_number before weight evaluation (#62148)
add 402d7885194 Java SDK: Resolve a task's arguments from the Dag's own
wiring (#73596)
add ca0d7328288 Move the Dag file processor's shared plumbing into a base
class
add 235be2d6bda Check that a serialized Dag can be stored and loaded
add bcba1cbf97f Fill a serialized Dag's unset settings from the Airflow
config
add 2f770c94c9a Reject malformed task entries in a serialized Dag
add fc73b8cc582 Check SDK Dags with the Dag processor's validation in
conformance
add 8cfd06019f4 Note in the parsing ADR that runtimes may omit
config-backed fields
add a63b9f1c091 Serialize the SDK Dag in tests that write a dag_maker Dag
add 0afbeb1e11c Sync the Java SDK copy of the Dag schema
add fd5eb37bd70 Let a coordinator parse the Dag files of the bundles it
serves
add 2cb5c904df8 Import the coordinator manager at the top of the importer
base
add 266131caf2a Check Dag file claims per bundle on the parse path
add 13d82cc57ec Drop serves_bundle so get_parsed_bundles is the only answer
add 6afe857c7e0 Say dag_bundle_name is only needed for two coordinators of
one class
add e610aa76488 Hand out a coordinator's Dag importer through
get_dag_importer alone
add 42c00d465ae Explain dag_bundle_name by the coordinator that parses a
Dag file
This update added new revisions after undoing existing revisions.
That is to say, some revisions that were in the old version of the
branch are not in the new version. This situation occurs
when a user --force pushes a change and generates a repository
containing something like this:
* -- * -- B -- O -- O -- O (d21aea99762)
\
N -- N -- N
refs/heads/jason/lang-sdk-e2e/03c-coordinator-dag-parsing (42c00d465ae)
You should already have received notification emails for all of the O
revisions, and so the following emails describe only the N revisions
from the common base, B.
Any revisions marked "omit" are not gone; other references still
refer to them. Any revisions marked "discard" are gone forever.
No new revisions were added by this update.
Summary of changes:
.../adr/lang-sdk/0010-native-dag-processing.md | 25 +-
.../authoring-and-scheduling/language-sdks/go.rst | 25 +-
.../language-sdks/java.rst | 10 +-
.../language-sdks/typescript.rst | 12 +-
airflow-core/src/airflow/models/taskinstance.py | 2 +-
.../ui/src/components/FilterBar/FilterBar.test.tsx | 21 +
.../ui/src/components/FilterBar/FilterBar.tsx | 31 +-
.../ui/src/pages/AssetsList/AssetsList.test.tsx | 85 +++-
.../airflow/ui/src/pages/AssetsList/AssetsList.tsx | 10 +-
airflow-core/tests/unit/models/test_dag.py | 2 +
go-sdk/README.md | 7 +-
java-sdk/README.md | 6 +-
.../org/apache/airflow/sdk/BuilderProcessor.kt | 2 +-
.../kotlin/org/apache/airflow/sdk/BuilderTest.kt | 4 +-
.../src/main/kotlin/org/apache/airflow/sdk/Arg.kt | 11 +-
.../main/kotlin/org/apache/airflow/sdk/Context.kt | 10 +-
.../main/kotlin/org/apache/airflow/sdk/DagDef.kt | 1 +
.../kotlin/org/apache/airflow/sdk/InputTask.kt | 2 +-
.../org/apache/airflow/sdk/execution/Task.kt | 7 +-
.../org/apache/airflow/sdk/internal/ArgValues.kt | 124 ++++++
.../org/apache/airflow/sdk/internal/TaskArgs.kt | 42 +-
.../org/apache/airflow/sdk/ArgTestSupport.kt | 15 +
.../kotlin/org/apache/airflow/sdk/ArgValuesTest.kt | 2 +-
.../kotlin/org/apache/airflow/sdk/InputTaskTest.kt | 90 ++++
.../org/apache/airflow/sdk/execution/TaskTest.kt | 24 +
.../apache/airflow/sdk/internal/ArgValuesTest.kt | 266 ++++++++++++
providers/apache/kafka/docs/triggers.rst | 33 ++
providers/apache/kafka/provider.yaml | 1 +
.../providers/apache/kafka/get_provider_info.py | 1 +
.../apache/kafka/triggers/shared_stream.py | 407 +++++++++++++++++
.../providers/apache/kafka/version_compat.py | 2 +
.../apache/kafka/triggers/test_shared_stream.py | 259 +++++++++++
.../apache/kafka/triggers/test_shared_stream.py | 481 +++++++++++++++++++++
.../src/airflow/sdk/coordinators/_subprocess.py | 11 +-
.../src/airflow/sdk/execution_time/coordinator.py | 93 ++--
task-sdk/src/airflow/sdk/importers/base.py | 18 +-
.../src/airflow/sdk/importers/python_importer.py | 4 +-
task-sdk/src/airflow/sdk/importers/zip_importer.py | 6 +-
.../tests/task_sdk/coordinators/test_subprocess.py | 12 -
.../task_sdk/execution_time/test_coordinator.py | 119 +++--
task-sdk/tests/task_sdk/importers/test_registry.py | 18 +-
ts-sdk/README.md | 8 +-
ts-sdk/example/README.md | 6 +-
43 files changed, 2050 insertions(+), 265 deletions(-)
create mode 100644
java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/internal/ArgValuesTest.kt
create mode 100644
providers/apache/kafka/src/airflow/providers/apache/kafka/triggers/shared_stream.py
create mode 100644
providers/apache/kafka/tests/integration/apache/kafka/triggers/test_shared_stream.py
create mode 100644
providers/apache/kafka/tests/unit/apache/kafka/triggers/test_shared_stream.py