guan404ming commented on code in PR #73441:
URL: https://github.com/apache/airflow/pull/73441#discussion_r4144776188


##########
ts-sdk/src/coordinator/serde.ts:
##########
@@ -0,0 +1,717 @@
+/*!
+ * 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.
+ */
+
+// Turns a Dag declared in TypeScript into Airflow's DagSerialization v3 JSON —
+// what the Dag processor stores and the scheduler reads. The format is
+// Airflow-internal rather than an SDK schema, so it is reimplemented per
+// language against `airflow-core/src/airflow/serialization/schema.json`; see
+// `airflow-core/adr/lang-sdk/0004-dag-parsing.md` for the field table this
+// follows, and the Java SDK's `Serde.kt` for the same job in another language.
+//
+// Byte-parity with Python's serializer is not the goal — Python omits fields
+// against a `client_defaults` table this SDK does not receive. What has to 
hold
+// is that `DagSerialization.from_dict` rebuilds the same Dag, which
+// scripts/ci/lang_sdk_serialization/compare.py checks against Python's own 
serialization.
+//
+// This module only produces the payload. Answering a Dag-parsing request with
+// it is the bundle's job, once the coordinator has a parse request to answer.
+
+import { relative as relativePath } from "node:path";
+
+import {
+  DAG_SCHEMA_FIELDS,
+  TASK_SCHEMA_FIELDS,
+  type SchemaField,
+} from "../generated/dag-schema-fields.js";
+import type { JsonValue } from "../sdk/client-types.js";
+import {
+  getDagOrderEdges,
+  getDagTaskGroups,
+  getDagTaskInputs,
+  getDagTaskRecords,
+  isTaskRef,
+  type Dag,
+  type RecordedInputs,
+  type TaskGroupRecord,
+} from "../sdk/dag.js";
+
+/** A serialized Dag: JSON, by the time it reaches the supervisor as msgpack. 
*/
+type SerializedValue = JsonValue;
+
+/** Airflow's type/var encoding, as `BaseSerialization.serialize()` emits it. 
*/
+interface TypeEncoded {
+  readonly __type: string;
+  readonly __var: SerializedValue;
+}
+
+/**
+ * Identity every TypeScript task carries, in place of the Python operator 
class
+ * a Python Dag would name.
+ *
+ * Fixed rather than derived: nothing on the Airflow side imports 
`_task_module`
+ * — `SerializedBaseOperator.populate_operator` only compares the pair as
+ * strings when matching plugin extra links — so the pair is free to name the
+ * coordinator that actually runs the task, which makes every TypeScript task
+ * greppable in the UI and the metadata DB.
+ */
+const TASK_TYPE = "TypeScriptOperator";
+const TASK_MODULE = "airflow.sdk.coordinators.node";
+
+/**
+ * Marks the tasks this SDK serialized, as the Java SDK marks its own.
+ *
+ * Nothing in airflow-core reads it today; the `operator` schema definition
+ * allows additional properties, so it rides along as a marker for tooling that
+ * wants to tell language-native tasks apart without parsing `_task_module`.
+ */
+const TASK_LANGUAGE = "typescript";
+
+/** How one set of authoring fields is written into a serialized object. */
+interface FieldRules {
+  readonly fields: Readonly<Record<string, SchemaField>>;
+  /**
+   * Fields whose {__type, __var} wrapper survives; every other field is
+   * unwrapped to the bare __var, as Python's `serialize_to_json` does.
+   *
+   * Neither set overlaps the generated authoring fields today, so in practice
+   * everything is unwrapped. They are named so that a decorated field added to
+   * the schema later takes the right path rather than silently losing its
+   * wrapper. Source: `DagSerialization._decorated_fields` and
+   * `OperatorSerialization._decorated_fields`.
+   */
+  readonly decorated: ReadonlySet<string>;
+  /** Fields never written, whatever the author set. */
+  readonly omitted: ReadonlySet<string>;
+}
+
+const DAG_FIELD_RULES: FieldRules = {
+  fields: DAG_SCHEMA_FIELDS,
+  decorated: new Set(["default_args", "access_control"]),
+  omitted: new Set(),
+};
+
+const TASK_FIELD_RULES: FieldRules = {
+  fields: TASK_SCHEMA_FIELDS,
+  decorated: new Set(["executor_config"]),
+  // Python drops both unless the operator names an email recipient
+  // (`OperatorSerialization._serialize_node`). A TypeScript task has no 
`email`
+  // field to name one, so writing them would describe a notification that can
+  // never be sent.
+  omitted: new Set(["email_on_failure", "email_on_retry"]),
+};
+
+const NULL_TIMETABLE = "airflow.timetables.simple.NullTimetable";
+const ONCE_TIMETABLE = "airflow.timetables.simple.OnceTimetable";
+const CONTINUOUS_TIMETABLE = "airflow.timetables.simple.ContinuousTimetable";
+const CRON_TIMETABLE = "airflow.timetables.trigger.CronTriggerTimetable";
+
+/** Serialize one Dag to the `dag` object of a DagSerialization v3 payload. */
+export function serializeDag(
+  dag: Dag,
+  fileloc: string,
+  relativeFileloc: string,
+): Record<string, SerializedValue> {
+  const graph = buildDagGraph(dag);
+  const inputs = getDagTaskInputs(dag);
+  const data: Record<string, SerializedValue> = {
+    dag_id: dag.dagId,
+    fileloc,
+    relative_fileloc: relativeFileloc,
+    timezone: "UTC",
+    timetable: serializeTimetable(dag.spec.schedule, dag.dagId),
+    tasks: [...getDagTaskRecords(dag)].map(([taskId, record]) =>
+      serializeTask(
+        dag.dagId,
+        taskId,
+        withDagQueue(record.spec, dag.spec.queue),
+        graph.downstreamTaskIds.get(taskId),
+        inputs.get(taskId),
+      ),
+    ),
+    dag_dependencies: [],
+    task_group: serializeTaskGroups(dag, graph),
+    edge_info: {},
+    params: [],
+    // Always written by Python's serializer, so a Dag without either still
+    // round-trips to the same object.
+    deadline: null,
+    allowed_run_types: null,
+  };
+  // A field the Dag leaves unset stays out, including the ones Python reads 
from
+  // config, such as max_active_tasks and catchup: Airflow fills those in from 
its
+  // own config when it receives the Dag.
+  applySchemaFields(data, dag.spec, DAG_FIELD_RULES, `Dag "${dag.dagId}"`);
+  return data;
+}
+
+/**
+ * The task's spec with the Dag's queue filled in, when the task named none.
+ *
+ * Merged before the fields are written rather than after, so a queue that
+ * happens to equal the schema default is omitted the way any other defaulted
+ * field is.
+ */
+function withDagQueue(spec: object, dagQueue: string | undefined): object {
+  if (dagQueue === undefined || "queue" in spec) return spec;
+  return { ...spec, queue: dagQueue };
+}
+
+/** Serialize one task, with its downstream edges sorted for a stable payload. 
*/
+function serializeTask(
+  dagId: string,
+  taskId: string,
+  spec: object,
+  downstream: ReadonlySet<string> | undefined,
+  inputs: RecordedInputs | undefined,
+): SerializedValue {
+  const data: Record<string, SerializedValue> = {
+    task_id: taskId,
+    task_type: TASK_TYPE,
+    _task_module: TASK_MODULE,
+    language: TASK_LANGUAGE,
+    // Python's operator serializer always emits this — its list value never
+    // matches the tuple default it is compared against. A TypeScript task has
+    // no Jinja templating, so the list is empty rather than absent.
+    template_fields: [],
+    // What marks a task whose arguments the API server resolves per instance
+    // and sends to a foreign runtime, as `@task.stub` does on the Python side.
+    // `get_arg_bindings` reads nothing without it.
+    is_stub: true,
+  };
+  const bindings = serializeArgBindings(inputs);
+  if (bindings) data["_arg_bindings"] = bindings;
+  applySchemaFields(data, spec, TASK_FIELD_RULES, `task "${taskId}" of Dag 
"${dagId}"`);
+  if (downstream?.size) {
+    data["downstream_task_ids"] = [...downstream].sort();
+  }
+  return { __type: "operator", __var: data };
+}
+
+/**
+ * The task's arguments as the binding spec the API server hands back at run
+ * time, one entry per argument in the order the call named them.
+ *
+ * A reference becomes an `xcom` binding naming the upstream task, and anything
+ * else a `literal` carrying the value, matching the `TaskArgBinding` union the
+ * execution API declares. `value_schema` is left out: it constrains the decode
+ * side, and nothing records a TypeScript argument's type at pack time, so
+ * omitting it says "unconstrained" rather than asserting a wrong type.
+ *
+ * `undefined` for a task called with no arguments, which needs no spec.
+ */
+function serializeArgBindings(inputs: RecordedInputs | undefined): 
SerializedValue | undefined {
+  if (inputs === undefined) return undefined;
+  const entries = Object.entries(inputs);
+  if (entries.length === 0) return undefined;
+  return entries.map(([name, value]): SerializedValue =>
+    isTaskRef(value)
+      ? { name, kind: "xcom", task_id: value.taskId }
+      : { name, kind: "literal", value: serializeValue(value) },

Review Comment:
   Literal `value` should stay plain JSON per the schema, but `serializeValue` 
wraps objects as `{__type, __var}`.



##########
scripts/ci/lang_sdk_serialization/compare.py:
##########
@@ -0,0 +1,233 @@
+# 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.
+r"""
+Check that a language SDK serializes Dags the way Airflow does.
+
+Runs the SDK's serializer and serialize_python.py over test_dags.yaml, each 
writing its serialization
+to a JSON file in a temporary directory. The Python side also takes the SDK's 
output as Airflow
+receives it, filling in the fields the SDK leaves to Airflow's config, and 
loads it through Airflow's
+deserializer. The two are then compared field by field, so this fails both 
when the serializers drift
+apart and when Airflow cannot read what the SDK writes.
+
+The SDK's serializer is the command after ``--``. It runs from the repository 
root, with the paths of
+test_dags.yaml and of the JSON file to write appended, and writes each Dag as
+``DagSerialization.to_dict`` would, keyed by Dag id. The Python side needs 
``uv``::
+
+    python3 scripts/ci/lang_sdk_serialization/compare.py --sdk typescript -- \
+        pnpm --dir ts-sdk exec tsx tests/conformance/serialize_typescript.ts
+
+On a failure the files are kept, and their directory is printed.
+"""
+
+from __future__ import annotations
+
+import argparse
+import json
+import shutil
+import subprocess
+import sys
+import tempfile
+from pathlib import Path
+from typing import Any
+
+HERE = Path(__file__).resolve().parent
+REPO_ROOT = HERE.parents[2]
+TEST_DAGS = HERE / "test_dags.yaml"
+SCHEMA = REPO_ROOT / "airflow-core" / "src" / "airflow" / "serialization" / 
"schema.json"
+
+# Name the file a Dag was declared in: Python's own Dag file, or the SDK's 
bundle.
+DAG_KEYS_NOT_COMPARED = frozenset({"fileloc", "relative_fileloc", 
"_processor_dags_folder"})
+
+TASK_KEYS_NOT_COMPARED = frozenset(
+    {
+        # Operator identity: Python names the Python class that ran, a 
language SDK a fixed pair.
+        # Neither is ever imported on the Airflow side.
+        "task_type",
+        "_task_module",
+        # The SDK's marker for a task in its language, which Python has no 
equivalent for.
+        "language",
+        # Python bookkeeping for mapped tasks and retry policies, neither of 
which a language SDK's
+        # Dag can declare.
+        "_needs_expansion",
+        "has_retry_policy",
+        # The stub contract every language SDK task carries, which the Python 
operator built here
+        # has no counterpart for. Loading the SDK's output through Airflow 
checks them.
+        "is_stub",
+        "_arg_bindings",

Review Comment:
   I think we could compare `_arg_bindings` too, so the check catches that 
wrapping.



-- 
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]

Reply via email to