jason810496 commented on code in PR #73441: URL: https://github.com/apache/airflow/pull/73441#discussion_r4144280105
########## ts-sdk/src/coordinator/serde.ts: ########## @@ -0,0 +1,729 @@ +/*! + * 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"; + +// Python resolves these from [core]/[scheduler] config when the Dag leaves them +// unset, and its serializer always writes the resolved value — there is no +// schema default to omit against. A bundle cannot read airflow.cfg, so the +// stock defaults stand in. +const DAG_CONFIG_FALLBACKS: Readonly<Record<string, SerializedValue>> = { + max_active_tasks: 16, // [core] max_active_tasks_per_dag + max_active_runs: 16, // [core] max_active_runs_per_dag + max_consecutive_failed_dag_runs: 0, + catchup: false, // [scheduler] catchup_by_default + disable_bundle_versioning: false, +}; + Review Comment: I think it would be better for the core Dag processing side handle it in #73842 (still WIP). Since all the SDKs will face this issue. Nice catch, thanks! -- 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]
