jason810496 commented on code in PR #73936: URL: https://github.com/apache/airflow/pull/73936#discussion_r4142615374
########## go-sdk/internal/genspec/authoring.go: ########## @@ -0,0 +1,445 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package main + +import ( + "encoding/json" + "errors" + "fmt" +) + +// authoringShape rewrites the two definitions the airflow package generates from +// into the shape a Dag author writes, rather than the shape Airflow serializes. +// Each definition drops the properties in exclude, rewrites the properties in +// override, and gains the properties in inject. +type authoringShape struct { + // doc becomes the description of the definition, which go-jsonschema writes as + // the doc comment of the generated type. + doc string + // exclude names each property that must not reach the generated struct, mapped + // to why. A property absent from the list generates, so that a property added + // on the Python side surfaces in review rather than vanishing; the reason is + // what a reviewer reads when deciding whether a new one belongs here. + exclude map[string]string + // override rewrites a property that generates as the wrong Go type. The + // serialization schema types a moment in time and a duration as a number of + // seconds and an integral count as a JSON number, none of which is the type an + // author sets. + override map[string]propertyOverride + // inject adds a property the schema has no counterpart for, so that every field + // of the generated struct comes from generation and the struct stays one + // declaration. + inject map[string]map[string]any +} + +// propertyOverride is the part of a property genspec rewrites. goType and imports +// become go-jsonschema's goJSONSchema extension, which it reads before a $ref, so +// an override applies to a property written as a reference too. +type propertyOverride struct { + goType string + imports []string + // doc replaces the description, and so the doc comment of the generated field, + // where the serialized property has nothing to say about how an author sets it. + doc string +} + +var authoringShapes = map[string]authoringShape{ + "dag": dagShape, + "operator": taskShape, +} + +var dagShape = authoringShape{ + doc: "DagSpec holds the attributes of a Dag other than its dag_id. Dag takes one.", + exclude: map[string]string{ + "dag_id": "a positional parameter of airflow.Dag, not a spec field", + "fileloc": "the path of the Dag file, which the bundle fills in", + "relative_fileloc": "the path of the Dag file, which the bundle fills in", + "_processor_dags_folder": "the Dag processor's own folder, filled in at parse time", + "bundle_name": "the name of the bundle that carries the Dag, not the Dag's", + "tasks": "the tasks dag.Task registers", + "task_group": "the groups dag.TaskGroup registers", + "edge_info": "the labels airflow.Label carries into an edge verb", + "dag_dependencies": "derived from the edges and the assets a Dag declares", + "timezone": "carried by the time.Time an author sets on StartDate", + "timetable": "the serialized form of Schedule, which is injected instead", + "allowed_run_types": "a union the author expresses by setting Schedule", + "_concurrency": "the pre-2.2 spelling of MaxActiveTasks", + "has_on_success_callback": "derived from whether a callback is registered", + "has_on_failure_callback": "derived from whether a callback is registered", + "params": "no Go authoring type yet: a param carries a schema of its own", + "default_args": "no Go authoring type yet: the values are arbitrary and untyped", + "access_control": "no Go authoring type yet, and it is deprecated in Airflow 3", + "owner_links": "no Go authoring type yet: an object of arbitrary link targets", + "deadline": "no Go authoring type yet: a serialized deadline reference", + "disable_bundle_versioning": "a property of the bundle, set where the bundle is configured", + "rerun_with_latest_version": "no Go authoring type yet: the tri-state a null allows", + }, + override: map[string]propertyOverride{ + "start_date": {goType: "time.Time", imports: []string{"time"}}, + "end_date": {goType: "time.Time", imports: []string{"time"}}, + "dagrun_timeout": {goType: "time.Duration", imports: []string{"time"}}, + "max_active_tasks": {goType: "int"}, + "max_active_runs": {goType: "int"}, + "max_consecutive_failed_dag_runs": {goType: "int"}, + "tags": {goType: "[]string"}, + }, + inject: map[string]map[string]any{ + // The schema carries the serialized timetable this resolves to, never the + // expression an author writes. + "schedule": { + "type": "string", + "description": "Schedule is the cron expression or preset the Dag runs on, such as \"@daily\".", + }, + }, +} + +var taskShape = authoringShape{ + doc: "TaskSpec holds the attributes of a task. DagRef.Task takes at most one per task.", + exclude: map[string]string{ + "task_type": "the operator class name, which the SDK fills in", + "_task_module": "the operator's Python module, which the SDK fills in", + "_operator_extra_links": "links a Python operator class declares, which a Go task has none of", + "ui_color": "the grid colour, which the SDK fills in", + "ui_fgcolor": "the grid colour, which the SDK fills in", + "template_fields": "the templated attributes of a Python operator class", + "template_ext": "the templated attributes of a Python operator class", + "template_fields_renderers": "the templated attributes of a Python operator class", + "downstream_task_ids": "the edges Before, After and Inputs declare", + "partial_kwargs": "the serialized form of a mapped task's partial arguments", + "_logger_name": "the logger the task runner names", + "_needs_expansion": "derived from whether the task is mapped", + "_is_mapped": "derived from whether the task is mapped", + "_is_sensor": "derived from the task's own kind", + "_disallow_kwargs_override": "a mapped-task serialization detail", + "_expand_input_attr": "a mapped-task serialization detail", + "_arg_bindings": "the bindings airflow.Inputs records", + "has_on_execute_callback": "derived from whether a callback is registered", + "has_on_failure_callback": "derived from whether a callback is registered", + "has_on_skipped_callback": "derived from whether a callback is registered", + "has_on_success_callback": "derived from whether a callback is registered", + "has_on_retry_callback": "derived from whether a callback is registered", + "start_from_trigger": "deferral is Python's, per decision 11 of ADR 8", + "start_trigger_args": "deferral is Python's, per decision 11 of ADR 8", + "multiple_outputs": "derived from the Go function's return type", + "params": "no Go authoring type yet: a param carries a schema of its own", + "executor_config": "no Go authoring type yet: the keys are executor-specific", + "inlets": "no Go authoring type yet: an asset needs its own spec", + "outlets": "no Go authoring type yet: an asset needs its own spec", + "render_template_as_native_obj": "set on the Dag, where the schema types it without a null", + "weight_rule": "needs a named type and its constants, as TriggerRule has", Review Comment: nit: `weight_rule` has a fixed set of values, so it can follow the `TriggerRule` pattern instead of being excluded, as #67155 exposes it: ```go type WeightRule string const ( WeightRuleDownstream WeightRule = "downstream" WeightRuleUpstream WeightRule = "upstream" WeightRuleAbsolute WeightRule = "absolute" ) ``` ```go "weight_rule": {goType: "WeightRule"}, ``` Plus a match-Python test like the trigger-rule one. ########## go-sdk/internal/genspec/authoring.go: ########## @@ -0,0 +1,445 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package main + +import ( + "encoding/json" + "errors" + "fmt" +) + +// authoringShape rewrites the two definitions the airflow package generates from +// into the shape a Dag author writes, rather than the shape Airflow serializes. +// Each definition drops the properties in exclude, rewrites the properties in +// override, and gains the properties in inject. +type authoringShape struct { + // doc becomes the description of the definition, which go-jsonschema writes as + // the doc comment of the generated type. + doc string + // exclude names each property that must not reach the generated struct, mapped + // to why. A property absent from the list generates, so that a property added + // on the Python side surfaces in review rather than vanishing; the reason is + // what a reviewer reads when deciding whether a new one belongs here. + exclude map[string]string + // override rewrites a property that generates as the wrong Go type. The + // serialization schema types a moment in time and a duration as a number of + // seconds and an integral count as a JSON number, none of which is the type an + // author sets. + override map[string]propertyOverride + // inject adds a property the schema has no counterpart for, so that every field + // of the generated struct comes from generation and the struct stays one + // declaration. + inject map[string]map[string]any +} + +// propertyOverride is the part of a property genspec rewrites. goType and imports +// become go-jsonschema's goJSONSchema extension, which it reads before a $ref, so +// an override applies to a property written as a reference too. +type propertyOverride struct { + goType string + imports []string + // doc replaces the description, and so the doc comment of the generated field, + // where the serialized property has nothing to say about how an author sets it. + doc string +} + +var authoringShapes = map[string]authoringShape{ + "dag": dagShape, + "operator": taskShape, +} + +var dagShape = authoringShape{ + doc: "DagSpec holds the attributes of a Dag other than its dag_id. Dag takes one.", + exclude: map[string]string{ + "dag_id": "a positional parameter of airflow.Dag, not a spec field", + "fileloc": "the path of the Dag file, which the bundle fills in", + "relative_fileloc": "the path of the Dag file, which the bundle fills in", + "_processor_dags_folder": "the Dag processor's own folder, filled in at parse time", + "bundle_name": "the name of the bundle that carries the Dag, not the Dag's", + "tasks": "the tasks dag.Task registers", + "task_group": "the groups dag.TaskGroup registers", + "edge_info": "the labels airflow.Label carries into an edge verb", + "dag_dependencies": "derived from the edges and the assets a Dag declares", + "timezone": "carried by the time.Time an author sets on StartDate", + "timetable": "the serialized form of Schedule, which is injected instead", + "allowed_run_types": "a union the author expresses by setting Schedule", + "_concurrency": "the pre-2.2 spelling of MaxActiveTasks", + "has_on_success_callback": "derived from whether a callback is registered", + "has_on_failure_callback": "derived from whether a callback is registered", + "params": "no Go authoring type yet: a param carries a schema of its own", + "default_args": "no Go authoring type yet: the values are arbitrary and untyped", + "access_control": "no Go authoring type yet, and it is deprecated in Airflow 3", + "owner_links": "no Go authoring type yet: an object of arbitrary link targets", + "deadline": "no Go authoring type yet: a serialized deadline reference", + "disable_bundle_versioning": "a property of the bundle, set where the bundle is configured", Review Comment: `disable_bundle_versioning` is a Dag argument (default from `[dag_processor]` config), not a bundle property. Without it, a Go Dag can't opt out of versioning. It has no schema default, so dropping the exclusion makes it `*bool`. ```suggestion ``` ########## go-sdk/airflow/spec.gen.go: ########## @@ -0,0 +1,184 @@ +// Code generated by github.com/atombender/go-jsonschema, DO NOT EDIT. +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package airflow + +import "time" + +// DagSpec holds the attributes of a Dag other than its dag_id. Dag takes one. +type DagSpec struct { + // Catchup corresponds to the JSON schema field "catchup". + Catchup *bool `json:"catchup,omitempty,omitzero"` + + // DagDisplayName corresponds to the JSON schema field "dag_display_name". + DagDisplayName string `json:"dag_display_name,omitempty,omitzero"` + + // DagrunTimeout corresponds to the JSON schema field "dagrun_timeout". + DagrunTimeout time.Duration `json:"dagrun_timeout,omitempty,omitzero"` + + // Description corresponds to the JSON schema field "description". + Description string `json:"description,omitempty,omitzero"` + + // DocMd corresponds to the JSON schema field "doc_md". + DocMd string `json:"doc_md,omitempty,omitzero"` + + // EndDate corresponds to the JSON schema field "end_date". + EndDate time.Time `json:"end_date,omitempty,omitzero"` + + // FailFast corresponds to the JSON schema field "fail_fast". + FailFast bool `json:"fail_fast,omitempty,omitzero"` + + // IsPausedUponCreation corresponds to the JSON schema field + // "is_paused_upon_creation". + IsPausedUponCreation *bool `json:"is_paused_upon_creation,omitempty,omitzero"` + + // MaxActiveRuns corresponds to the JSON schema field "max_active_runs". + MaxActiveRuns int `json:"max_active_runs,omitempty,omitzero"` + + // MaxActiveTasks corresponds to the JSON schema field "max_active_tasks". + MaxActiveTasks int `json:"max_active_tasks,omitempty,omitzero"` + + // MaxConsecutiveFailedDagRuns corresponds to the JSON schema field + // "max_consecutive_failed_dag_runs". + MaxConsecutiveFailedDagRuns int `json:"max_consecutive_failed_dag_runs,omitempty,omitzero"` + + // RenderTemplateAsNativeObj corresponds to the JSON schema field + // "render_template_as_native_obj". + RenderTemplateAsNativeObj bool `json:"render_template_as_native_obj,omitempty,omitzero"` + + // Schedule is the cron expression or preset the Dag runs on, such as "@daily". + Schedule string `json:"schedule,omitempty,omitzero"` + + // StartDate corresponds to the JSON schema field "start_date". + StartDate time.Time `json:"start_date,omitempty,omitzero"` + + // Tags corresponds to the JSON schema field "tags". + Tags []string `json:"tags,omitempty,omitzero"` +} + +// TaskSpec holds the attributes of a task. DagRef.Task takes at most one per task. +type TaskSpec struct { + // TaskDisplayName corresponds to the JSON schema field "_task_display_name". + TaskDisplayName string `json:"_task_display_name,omitempty,omitzero"` + + // AllowNestedOperators corresponds to the JSON schema field + // "allow_nested_operators". + AllowNestedOperators *bool `json:"allow_nested_operators,omitempty,omitzero"` Review Comment: `allow_nested_operators` is Python-only, and the legacy `doc_*` fields shouldn't become public API. #67155 keeps only `doc_md`. Adding fields later is easy, removing them isn't: ```go "allow_nested_operators": "Python-only: warns when an operator executes inside another operator", "doc": "legacy; only doc_md is exposed", "doc_json": "legacy; only doc_md is exposed", "doc_rst": "legacy; only doc_md is exposed", "doc_yaml": "legacy; only doc_md is exposed", ``` ########## go-sdk/airflow/spec.go: ########## @@ -17,12 +17,38 @@ package airflow -// DagSpec holds the attributes of a Dag other than its dag_id. [Dag] takes one. -type DagSpec struct{} +// DagSpec and TaskSpec are generated from Airflow core's Dag serialization schema, +// which Python owns, so that neither struct drifts from it silently. genspec +// rewrites the schema into the authoring shape, go-jsonschema writes the structs, +// and genspec puts the license header back on what it wrote. The rewritten schema +// is a build artifact under .build; spec.gen.go is committed. +// +// To change a field, change the schema on the Python side, or the exclusions, type +// overrides and injected properties in internal/genspec/authoring.go, and run +// `just generate-specs`. + +//go:generate go run ../internal/genspec -schema ../../airflow-core/src/airflow/serialization/schema.json -out ../../.build/go-sdk/spec.schema.json +//go:generate go run github.com/atombender/[email protected] --only-models --struct-name-from-title --tags json --capitalization ID --capitalization JSON -p airflow -o spec.gen.go ../../.build/go-sdk/spec.schema.json Review Comment: These `json` tags aren't the wire format: `encoding/json` writes durations as nanoseconds, and `omitempty` drops explicit `false`/`0`. A later `json.Marshal(spec)` would produce a shape core misreads. Drop the tags, or document that they aren't the wire format and port #67155's `SchemaFields()` in the PR. ########## go-sdk/internal/genspec/authoring.go: ########## @@ -0,0 +1,445 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package main + +import ( + "encoding/json" + "errors" + "fmt" +) + +// authoringShape rewrites the two definitions the airflow package generates from +// into the shape a Dag author writes, rather than the shape Airflow serializes. +// Each definition drops the properties in exclude, rewrites the properties in +// override, and gains the properties in inject. +type authoringShape struct { + // doc becomes the description of the definition, which go-jsonschema writes as + // the doc comment of the generated type. + doc string + // exclude names each property that must not reach the generated struct, mapped + // to why. A property absent from the list generates, so that a property added + // on the Python side surfaces in review rather than vanishing; the reason is + // what a reviewer reads when deciding whether a new one belongs here. + exclude map[string]string + // override rewrites a property that generates as the wrong Go type. The + // serialization schema types a moment in time and a duration as a number of + // seconds and an integral count as a JSON number, none of which is the type an + // author sets. + override map[string]propertyOverride + // inject adds a property the schema has no counterpart for, so that every field + // of the generated struct comes from generation and the struct stays one + // declaration. + inject map[string]map[string]any +} + +// propertyOverride is the part of a property genspec rewrites. goType and imports +// become go-jsonschema's goJSONSchema extension, which it reads before a $ref, so +// an override applies to a property written as a reference too. +type propertyOverride struct { + goType string + imports []string + // doc replaces the description, and so the doc comment of the generated field, + // where the serialized property has nothing to say about how an author sets it. + doc string +} + +var authoringShapes = map[string]authoringShape{ + "dag": dagShape, + "operator": taskShape, +} + +var dagShape = authoringShape{ + doc: "DagSpec holds the attributes of a Dag other than its dag_id. Dag takes one.", + exclude: map[string]string{ + "dag_id": "a positional parameter of airflow.Dag, not a spec field", + "fileloc": "the path of the Dag file, which the bundle fills in", + "relative_fileloc": "the path of the Dag file, which the bundle fills in", + "_processor_dags_folder": "the Dag processor's own folder, filled in at parse time", + "bundle_name": "the name of the bundle that carries the Dag, not the Dag's", + "tasks": "the tasks dag.Task registers", + "task_group": "the groups dag.TaskGroup registers", + "edge_info": "the labels airflow.Label carries into an edge verb", + "dag_dependencies": "derived from the edges and the assets a Dag declares", + "timezone": "carried by the time.Time an author sets on StartDate", + "timetable": "the serialized form of Schedule, which is injected instead", + "allowed_run_types": "a union the author expresses by setting Schedule", + "_concurrency": "the pre-2.2 spelling of MaxActiveTasks", + "has_on_success_callback": "derived from whether a callback is registered", + "has_on_failure_callback": "derived from whether a callback is registered", + "params": "no Go authoring type yet: a param carries a schema of its own", + "default_args": "no Go authoring type yet: the values are arbitrary and untyped", + "access_control": "no Go authoring type yet, and it is deprecated in Airflow 3", + "owner_links": "no Go authoring type yet: an object of arbitrary link targets", + "deadline": "no Go authoring type yet: a serialized deadline reference", + "disable_bundle_versioning": "a property of the bundle, set where the bundle is configured", + "rerun_with_latest_version": "no Go authoring type yet: the tri-state a null allows", + }, + override: map[string]propertyOverride{ + "start_date": {goType: "time.Time", imports: []string{"time"}}, + "end_date": {goType: "time.Time", imports: []string{"time"}}, + "dagrun_timeout": {goType: "time.Duration", imports: []string{"time"}}, + "max_active_tasks": {goType: "int"}, + "max_active_runs": {goType: "int"}, + "max_consecutive_failed_dag_runs": {goType: "int"}, + "tags": {goType: "[]string"}, + }, + inject: map[string]map[string]any{ + // The schema carries the serialized timetable this resolves to, never the + // expression an author writes. + "schedule": { + "type": "string", + "description": "Schedule is the cron expression or preset the Dag runs on, such as \"@daily\".", + }, + }, +} + +var taskShape = authoringShape{ + doc: "TaskSpec holds the attributes of a task. DagRef.Task takes at most one per task.", + exclude: map[string]string{ + "task_type": "the operator class name, which the SDK fills in", + "_task_module": "the operator's Python module, which the SDK fills in", + "_operator_extra_links": "links a Python operator class declares, which a Go task has none of", + "ui_color": "the grid colour, which the SDK fills in", + "ui_fgcolor": "the grid colour, which the SDK fills in", + "template_fields": "the templated attributes of a Python operator class", + "template_ext": "the templated attributes of a Python operator class", + "template_fields_renderers": "the templated attributes of a Python operator class", + "downstream_task_ids": "the edges Before, After and Inputs declare", + "partial_kwargs": "the serialized form of a mapped task's partial arguments", + "_logger_name": "the logger the task runner names", + "_needs_expansion": "derived from whether the task is mapped", + "_is_mapped": "derived from whether the task is mapped", + "_is_sensor": "derived from the task's own kind", + "_disallow_kwargs_override": "a mapped-task serialization detail", + "_expand_input_attr": "a mapped-task serialization detail", + "_arg_bindings": "the bindings airflow.Inputs records", + "has_on_execute_callback": "derived from whether a callback is registered", + "has_on_failure_callback": "derived from whether a callback is registered", + "has_on_skipped_callback": "derived from whether a callback is registered", + "has_on_success_callback": "derived from whether a callback is registered", + "has_on_retry_callback": "derived from whether a callback is registered", + "start_from_trigger": "deferral is Python's, per decision 11 of ADR 8", + "start_trigger_args": "deferral is Python's, per decision 11 of ADR 8", + "multiple_outputs": "derived from the Go function's return type", + "params": "no Go authoring type yet: a param carries a schema of its own", + "executor_config": "no Go authoring type yet: the keys are executor-specific", + "inlets": "no Go authoring type yet: an asset needs its own spec", + "outlets": "no Go authoring type yet: an asset needs its own spec", + "render_template_as_native_obj": "set on the Dag, where the schema types it without a null", + "weight_rule": "needs a named type and its constants, as TriggerRule has", + // Airflow's BaseOperator attribute is a bool, so generating the number the + // schema declares would put a float64 where an author writes true. + "retry_exponential_backoff": "the schema types it as a number; the attribute is a bool", Review Comment: `retry_exponential_backoff` is a `float` multiplier, not a bool, so Go authors can't enable backoff. #67155 exposes it as `float64`. ```suggestion ``` and in `taskShape.override`: ```go "retry_exponential_backoff": {goType: "float64"}, ``` ########## go-sdk/airflow/spec.gen.go: ########## @@ -0,0 +1,184 @@ +// Code generated by github.com/atombender/go-jsonschema, DO NOT EDIT. +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package airflow + +import "time" + +// DagSpec holds the attributes of a Dag other than its dag_id. Dag takes one. +type DagSpec struct { + // Catchup corresponds to the JSON schema field "catchup". + Catchup *bool `json:"catchup,omitempty,omitzero"` + + // DagDisplayName corresponds to the JSON schema field "dag_display_name". + DagDisplayName string `json:"dag_display_name,omitempty,omitzero"` + + // DagrunTimeout corresponds to the JSON schema field "dagrun_timeout". + DagrunTimeout time.Duration `json:"dagrun_timeout,omitempty,omitzero"` + + // Description corresponds to the JSON schema field "description". + Description string `json:"description,omitempty,omitzero"` + + // DocMd corresponds to the JSON schema field "doc_md". + DocMd string `json:"doc_md,omitempty,omitzero"` + + // EndDate corresponds to the JSON schema field "end_date". + EndDate time.Time `json:"end_date,omitempty,omitzero"` + + // FailFast corresponds to the JSON schema field "fail_fast". + FailFast bool `json:"fail_fast,omitempty,omitzero"` + + // IsPausedUponCreation corresponds to the JSON schema field + // "is_paused_upon_creation". + IsPausedUponCreation *bool `json:"is_paused_upon_creation,omitempty,omitzero"` + + // MaxActiveRuns corresponds to the JSON schema field "max_active_runs". + MaxActiveRuns int `json:"max_active_runs,omitempty,omitzero"` + + // MaxActiveTasks corresponds to the JSON schema field "max_active_tasks". + MaxActiveTasks int `json:"max_active_tasks,omitempty,omitzero"` + + // MaxConsecutiveFailedDagRuns corresponds to the JSON schema field + // "max_consecutive_failed_dag_runs". + MaxConsecutiveFailedDagRuns int `json:"max_consecutive_failed_dag_runs,omitempty,omitzero"` + + // RenderTemplateAsNativeObj corresponds to the JSON schema field + // "render_template_as_native_obj". + RenderTemplateAsNativeObj bool `json:"render_template_as_native_obj,omitempty,omitzero"` + + // Schedule is the cron expression or preset the Dag runs on, such as "@daily". + Schedule string `json:"schedule,omitempty,omitzero"` + + // StartDate corresponds to the JSON schema field "start_date". + StartDate time.Time `json:"start_date,omitempty,omitzero"` + + // Tags corresponds to the JSON schema field "tags". + Tags []string `json:"tags,omitempty,omitzero"` +} + +// TaskSpec holds the attributes of a task. DagRef.Task takes at most one per task. +type TaskSpec struct { + // TaskDisplayName corresponds to the JSON schema field "_task_display_name". + TaskDisplayName string `json:"_task_display_name,omitempty,omitzero"` + + // AllowNestedOperators corresponds to the JSON schema field + // "allow_nested_operators". + AllowNestedOperators *bool `json:"allow_nested_operators,omitempty,omitzero"` + + // DependsOnPast corresponds to the JSON schema field "depends_on_past". + DependsOnPast bool `json:"depends_on_past,omitempty,omitzero"` + + // DoXcomPush corresponds to the JSON schema field "do_xcom_push". + DoXcomPush *bool `json:"do_xcom_push,omitempty,omitzero"` + + // Doc corresponds to the JSON schema field "doc". + Doc string `json:"doc,omitempty,omitzero"` + + // DocJSON corresponds to the JSON schema field "doc_json". + DocJSON string `json:"doc_json,omitempty,omitzero"` + + // DocMd corresponds to the JSON schema field "doc_md". + DocMd string `json:"doc_md,omitempty,omitzero"` + + // DocRst corresponds to the JSON schema field "doc_rst". + DocRst string `json:"doc_rst,omitempty,omitzero"` + + // DocYaml corresponds to the JSON schema field "doc_yaml". + DocYaml string `json:"doc_yaml,omitempty,omitzero"` + + // EmailOnFailure corresponds to the JSON schema field "email_on_failure". + EmailOnFailure *bool `json:"email_on_failure,omitempty,omitzero"` + + // EmailOnRetry corresponds to the JSON schema field "email_on_retry". + EmailOnRetry *bool `json:"email_on_retry,omitempty,omitzero"` + + // EndDate corresponds to the JSON schema field "end_date". + EndDate time.Time `json:"end_date,omitempty,omitzero"` + + // ExecutionTimeout corresponds to the JSON schema field "execution_timeout". + ExecutionTimeout time.Duration `json:"execution_timeout,omitempty,omitzero"` + + // Executor corresponds to the JSON schema field "executor". + Executor string `json:"executor,omitempty,omitzero"` + + // IgnoreFirstDependsOnPast corresponds to the JSON schema field + // "ignore_first_depends_on_past". + IgnoreFirstDependsOnPast bool `json:"ignore_first_depends_on_past,omitempty,omitzero"` + + // IsSetup corresponds to the JSON schema field "is_setup". + IsSetup bool `json:"is_setup,omitempty,omitzero"` Review Comment: `TaskSpec{IsTeardown: true}` keeps `trigger_rule=all_success`, so the teardown is skipped when upstream fails. Python's `as_teardown()` also sets `all_done_setup_success`. Exclude these until the SDK models setup/teardown, as #67155 does: ```go "is_setup": "setup/teardown needs trigger-rule handling the SDK does not model yet", "is_teardown": "setup/teardown needs trigger-rule handling the SDK does not model yet", "on_failure_fail_dagrun": "only meaningful on a teardown task", ``` ########## go-sdk/airflow/spec.go: ########## @@ -17,12 +17,38 @@ package airflow -// DagSpec holds the attributes of a Dag other than its dag_id. [Dag] takes one. -type DagSpec struct{} +// DagSpec and TaskSpec are generated from Airflow core's Dag serialization schema, +// which Python owns, so that neither struct drifts from it silently. genspec +// rewrites the schema into the authoring shape, go-jsonschema writes the structs, +// and genspec puts the license header back on what it wrote. The rewritten schema +// is a build artifact under .build; spec.gen.go is committed. +// +// To change a field, change the schema on the Python side, or the exclusions, type +// overrides and injected properties in internal/genspec/authoring.go, and run +// `just generate-specs`. + +//go:generate go run ../internal/genspec -schema ../../airflow-core/src/airflow/serialization/schema.json -out ../../.build/go-sdk/spec.schema.json +//go:generate go run github.com/atombender/[email protected] --only-models --struct-name-from-title --tags json --capitalization ID --capitalization JSON -p airflow -o spec.gen.go ../../.build/go-sdk/spec.schema.json +//go:generate go run ../internal/genspec -license spec.gen.go + +// TriggerRule is when a task runs, given the state of the tasks upstream of it. +// The serialization schema types trigger_rule as a plain string and names none of +// its values, so the constants are written here and the generated field is given +// this type by internal/genspec/authoring.go. +type TriggerRule string Review Comment: For this one, let's track a issue for now instead of directly bumping to Go 1.26 (which is kind of "too new" for adoption, as it just released in Feb this year). Another direction that Claude suggested is to have a `Ptr` generic function helper, but I felt it's not Airflow's responsibility to support the pointer transformation. --- nit: several fields are now pointers, and `&false` doesn't compile. Go 1.26 `new(expr)` covers this without a helper, so bump `go.mod` to `go 1.26` and use it in docs and examples: ```go airflow.DagSpec{Catchup: new(false)} ``` ########## go-sdk/internal/genspec/authoring.go: ########## @@ -0,0 +1,445 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package main + +import ( + "encoding/json" + "errors" + "fmt" +) + +// authoringShape rewrites the two definitions the airflow package generates from +// into the shape a Dag author writes, rather than the shape Airflow serializes. +// Each definition drops the properties in exclude, rewrites the properties in +// override, and gains the properties in inject. +type authoringShape struct { + // doc becomes the description of the definition, which go-jsonschema writes as + // the doc comment of the generated type. + doc string + // exclude names each property that must not reach the generated struct, mapped + // to why. A property absent from the list generates, so that a property added + // on the Python side surfaces in review rather than vanishing; the reason is + // what a reviewer reads when deciding whether a new one belongs here. + exclude map[string]string + // override rewrites a property that generates as the wrong Go type. The + // serialization schema types a moment in time and a duration as a number of + // seconds and an integral count as a JSON number, none of which is the type an + // author sets. + override map[string]propertyOverride + // inject adds a property the schema has no counterpart for, so that every field + // of the generated struct comes from generation and the struct stays one + // declaration. + inject map[string]map[string]any +} + +// propertyOverride is the part of a property genspec rewrites. goType and imports +// become go-jsonschema's goJSONSchema extension, which it reads before a $ref, so +// an override applies to a property written as a reference too. +type propertyOverride struct { + goType string + imports []string + // doc replaces the description, and so the doc comment of the generated field, + // where the serialized property has nothing to say about how an author sets it. + doc string +} + +var authoringShapes = map[string]authoringShape{ + "dag": dagShape, + "operator": taskShape, +} + +var dagShape = authoringShape{ + doc: "DagSpec holds the attributes of a Dag other than its dag_id. Dag takes one.", + exclude: map[string]string{ + "dag_id": "a positional parameter of airflow.Dag, not a spec field", + "fileloc": "the path of the Dag file, which the bundle fills in", + "relative_fileloc": "the path of the Dag file, which the bundle fills in", + "_processor_dags_folder": "the Dag processor's own folder, filled in at parse time", + "bundle_name": "the name of the bundle that carries the Dag, not the Dag's", + "tasks": "the tasks dag.Task registers", + "task_group": "the groups dag.TaskGroup registers", + "edge_info": "the labels airflow.Label carries into an edge verb", + "dag_dependencies": "derived from the edges and the assets a Dag declares", + "timezone": "carried by the time.Time an author sets on StartDate", + "timetable": "the serialized form of Schedule, which is injected instead", + "allowed_run_types": "a union the author expresses by setting Schedule", + "_concurrency": "the pre-2.2 spelling of MaxActiveTasks", + "has_on_success_callback": "derived from whether a callback is registered", + "has_on_failure_callback": "derived from whether a callback is registered", + "params": "no Go authoring type yet: a param carries a schema of its own", + "default_args": "no Go authoring type yet: the values are arbitrary and untyped", + "access_control": "no Go authoring type yet, and it is deprecated in Airflow 3", + "owner_links": "no Go authoring type yet: an object of arbitrary link targets", + "deadline": "no Go authoring type yet: a serialized deadline reference", + "disable_bundle_versioning": "a property of the bundle, set where the bundle is configured", + "rerun_with_latest_version": "no Go authoring type yet: the tri-state a null allows", + }, + override: map[string]propertyOverride{ + "start_date": {goType: "time.Time", imports: []string{"time"}}, + "end_date": {goType: "time.Time", imports: []string{"time"}}, + "dagrun_timeout": {goType: "time.Duration", imports: []string{"time"}}, + "max_active_tasks": {goType: "int"}, + "max_active_runs": {goType: "int"}, + "max_consecutive_failed_dag_runs": {goType: "int"}, + "tags": {goType: "[]string"}, + }, + inject: map[string]map[string]any{ + // The schema carries the serialized timetable this resolves to, never the + // expression an author writes. + "schedule": { + "type": "string", + "description": "Schedule is the cron expression or preset the Dag runs on, such as \"@daily\".", + }, + }, +} + +var taskShape = authoringShape{ + doc: "TaskSpec holds the attributes of a task. DagRef.Task takes at most one per task.", + exclude: map[string]string{ + "task_type": "the operator class name, which the SDK fills in", + "_task_module": "the operator's Python module, which the SDK fills in", + "_operator_extra_links": "links a Python operator class declares, which a Go task has none of", + "ui_color": "the grid colour, which the SDK fills in", + "ui_fgcolor": "the grid colour, which the SDK fills in", + "template_fields": "the templated attributes of a Python operator class", + "template_ext": "the templated attributes of a Python operator class", + "template_fields_renderers": "the templated attributes of a Python operator class", + "downstream_task_ids": "the edges Before, After and Inputs declare", + "partial_kwargs": "the serialized form of a mapped task's partial arguments", + "_logger_name": "the logger the task runner names", + "_needs_expansion": "derived from whether the task is mapped", + "_is_mapped": "derived from whether the task is mapped", + "_is_sensor": "derived from the task's own kind", + "_disallow_kwargs_override": "a mapped-task serialization detail", + "_expand_input_attr": "a mapped-task serialization detail", + "_arg_bindings": "the bindings airflow.Inputs records", + "has_on_execute_callback": "derived from whether a callback is registered", + "has_on_failure_callback": "derived from whether a callback is registered", + "has_on_skipped_callback": "derived from whether a callback is registered", + "has_on_success_callback": "derived from whether a callback is registered", + "has_on_retry_callback": "derived from whether a callback is registered", + "start_from_trigger": "deferral is Python's, per decision 11 of ADR 8", + "start_trigger_args": "deferral is Python's, per decision 11 of ADR 8", + "multiple_outputs": "derived from the Go function's return type", + "params": "no Go authoring type yet: a param carries a schema of its own", + "executor_config": "no Go authoring type yet: the keys are executor-specific", + "inlets": "no Go authoring type yet: an asset needs its own spec", + "outlets": "no Go authoring type yet: an asset needs its own spec", + "render_template_as_native_obj": "set on the Dag, where the schema types it without a null", + "weight_rule": "needs a named type and its constants, as TriggerRule has", + // Airflow's BaseOperator attribute is a bool, so generating the number the + // schema declares would put a float64 where an author writes true. + "retry_exponential_backoff": "the schema types it as a number; the attribute is a bool", + }, + override: map[string]propertyOverride{ + "start_date": {goType: "time.Time", imports: []string{"time"}}, + "end_date": {goType: "time.Time", imports: []string{"time"}}, + "execution_timeout": {goType: "time.Duration", imports: []string{"time"}}, + "retry_delay": {goType: "time.Duration", imports: []string{"time"}}, + "max_retry_delay": {goType: "time.Duration", imports: []string{"time"}}, + "retries": {goType: "int"}, + "pool_slots": {goType: "int"}, + "priority_weight": {goType: "int"}, + "max_active_tis_per_dag": {goType: "int"}, + "max_active_tis_per_dagrun": {goType: "int"}, + // TriggerRule and its constants are hand-written in the airflow package: the + // schema types the field as a plain string and names none of its values. + "trigger_rule": {goType: "TriggerRule"}, + "task_id": { + goType: "string", + doc: "TaskID is the task_id of the task. When TaskID is empty, the task_id is the name of the Go function that the task runs.", + }, + }, +} + +// shapeForAuthoring rewrites doc into the schema the spec structs generate from: +// each definition in shapes takes its authoring shape, and everything the airflow +// package does not generate is dropped. +func shapeForAuthoring(doc map[string]any, shapes map[string]authoringShape) error { + if err := applyAuthoringShapes(doc, shapes); err != nil { + return err + } + return keepOnlySpecGeneratingSchema(doc, shapes) +} + +// applyAuthoringShapes rewrites each definition in shapes. It reports a list entry +// that no longer matches the schema — an excluded or overridden property that has +// gone, an injected property the schema has grown — because each of those means the +// list here decides nothing and the generated struct would silently change shape. +func applyAuthoringShapes(doc map[string]any, shapes map[string]authoringShape) error { + definitions, ok := doc["definitions"].(map[string]any) + if !ok { + return errNoDefinitions + } + for _, name := range sortedKeys(shapes) { + definition, ok := definitions[name].(map[string]any) + if !ok { + return fmt.Errorf( + "definitions/%s is missing, and the airflow package generates from it", + name, + ) + } + properties, ok := definition["properties"].(map[string]any) + if !ok { + return fmt.Errorf("definitions/%s has no properties to generate a struct from", name) + } + shape := shapes[name] + if err := excludeProperties(name, definition, properties, shape.exclude); err != nil { + return err + } + if err := overrideProperties(name, properties, shape.override); err != nil { + return err + } + if err := injectProperties(name, properties, shape.inject); err != nil { + return err + } + if err := rejectCombinators(name, definition); err != nil { + return err + } + definition["description"] = shape.doc + // An authoring struct holds the fields it declares and no others, and every + // field is optional: airflow.Dag takes the dag_id positionally and a task_id + // defaults to the name of the Go function. + definition["additionalProperties"] = false + delete(definition, "required") + setPointers(properties) + } + return nil +} + +func excludeProperties( + name string, definition, properties map[string]any, exclude map[string]string, +) error { + for _, property := range sortedKeys(exclude) { + if _, ok := properties[property]; !ok { + return fmt.Errorf( + "definitions/%s/properties/%s is excluded as %q, but the schema no longer has it", + name, property, exclude[property], + ) + } + delete(properties, property) + } + return nil +} + +// overrideProperties replaces the Go type of a property with goJSONSchema, the +// extension go-jsonschema reads before it reads a type or a $ref. +func overrideProperties( + name string, + properties map[string]any, + override map[string]propertyOverride, +) error { + for _, property := range sortedKeys(override) { + node, ok := properties[property].(map[string]any) + if !ok { + return fmt.Errorf( + "definitions/%s/properties/%s is overridden to %s, but the schema no longer has it", + name, property, override[property].goType, + ) + } + extension := map[string]any{"type": override[property].goType} + if imports := override[property].imports; len(imports) > 0 { + extension["imports"] = anySlice(imports) + } + if doc := override[property].doc; doc != "" { + node["description"] = doc + } + node["goJSONSchema"] = extension + // The override replaces whatever the reference resolves to, and dropping it + // keeps the pruned schema free of references to definitions that are gone. + delete(node, "$ref") + } + return nil +} + +// rejectCombinators reports an anyOf, oneOf or allOf left in a definition the specs +// generate from. genspec has no rule for one, and go-jsonschema does not fail on it +// either: it degrades the property to interface{} or gives it a typedef of its own, +// so the field's meaning is lost in a diff that still looks like a field. The +// nullable pair the schema writes as anyOf [boolean, null] is the shape to expect, +// and resolveNullableTypes only handles the type-list spelling of it. +// +// A property whose type an override replaces is exempt: the override stands for +// whatever the schema says the property is. +func rejectCombinators(name string, definition map[string]any) error { + return walkFrom("/definitions/"+name, definition, func(path string, node map[string]any) error { + if extension, ok := node["goJSONSchema"].(map[string]any); ok { + if _, overridden := extension["type"]; overridden { + return nil + } + } + for _, keyword := range []string{"allOf", "anyOf", "oneOf"} { + if _, ok := node[keyword]; ok { + return fmt.Errorf( + "%s has %s, which genspec has no rule for; exclude the property or give it "+ + "a type override", + path, keyword, + ) + } + } + return nil + }) +} + +// setPointers decides, for every property left, whether its field is a pointer. A +// field needs one wherever the Go zero value is something an author could mean and +// the schema does not assert that it is already the default: a concrete field would +// make that setting indistinguishable from an unset one, and omitempty would drop +// it on the way out. +// +// Two shapes qualify. A scalar whose schema default is not the Go zero value is the +// rule pkg/execution/genmodels applies to the supervisor schema. A boolean with no +// schema default at all is the second: false is always one of its two legal values, +// and the absence of a default does not mean the default is false — catchup and +// is_paused_upon_creation take theirs from [scheduler] catchup_by_default and +// [core] dags_are_paused_at_creation, both of which can be true. A count keeps its +// concrete type under the same reasoning: 0 is not a value max_active_runs or +// max_active_tis_per_dag can take, so the zero value can only mean unset. Review Comment: Not true for `max_consecutive_failed_dag_runs`: `0` means "never auto-pause", and the default comes from config, which can be non-zero. As a plain `int`, an author can't set `0` explicitly. Let an override force a pointer: ```go type propertyOverride struct { goType string imports []string doc string pointer bool } ``` ```go "max_consecutive_failed_dag_runs": {goType: "int", pointer: true}, ``` ########## go-sdk/airflow/dag.go: ########## @@ -64,6 +65,10 @@ func Dag(dagID string, spec ...DagSpec) *DagRef { d := &DagRef{dagID: dagID} if len(spec) == 1 { d.spec = spec[0] + // Assigning the struct copies the slice header, not what it points at. + // TestDagAndTaskCopyTheSliceAndMapFieldsOfTheirSpecs fails when a spec grows another + // slice or map field and it is not copied here. + d.spec.Tags = slices.Clone(spec[0].Tags) Review Comment: Pointer fields are shared with the caller's spec, so the registered Dag can still change: ```go catchup := true d := airflow.Dag("etl", airflow.DagSpec{Catchup: &catchup}) catchup = false // registered Dag now has catchup=false ``` The copy test skips `reflect.Pointer` fields. Copy pointers in `Dag`/`Task` too and cover them in the test. ########## go-sdk/airflow/spec.go: ########## @@ -17,12 +17,38 @@ package airflow -// DagSpec holds the attributes of a Dag other than its dag_id. [Dag] takes one. -type DagSpec struct{} +// DagSpec and TaskSpec are generated from Airflow core's Dag serialization schema, +// which Python owns, so that neither struct drifts from it silently. genspec +// rewrites the schema into the authoring shape, go-jsonschema writes the structs, +// and genspec puts the license header back on what it wrote. The rewritten schema +// is a build artifact under .build; spec.gen.go is committed. +// +// To change a field, change the schema on the Python side, or the exclusions, type +// overrides and injected properties in internal/genspec/authoring.go, and run +// `just generate-specs`. + +//go:generate go run ../internal/genspec -schema ../../airflow-core/src/airflow/serialization/schema.json -out ../../.build/go-sdk/spec.schema.json +//go:generate go run github.com/atombender/[email protected] --only-models --struct-name-from-title --tags json --capitalization ID --capitalization JSON -p airflow -o spec.gen.go ../../.build/go-sdk/spec.schema.json Review Comment: nit: `DocMd` and `DoXcomPush` should be `DocMD` and `DoXComPush`. Two more capitalizations fix it: ```suggestion //go:generate go run github.com/atombender/[email protected] --only-models --struct-name-from-title --tags json --capitalization ID --capitalization JSON --capitalization MD --capitalization XCom -p airflow -o spec.gen.go ../../.build/go-sdk/spec.schema.json ``` -- 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]
