uranusjr commented on code in PR #71929: URL: https://github.com/apache/airflow/pull/71929#discussion_r4129588344
########## airflow-core/adr/lang-sdk/0012-lang-sdk-parse-protocol.md: ########## @@ -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. + --> + +# ADR-0012: Lang-SDK Parse Protocol — Handler Messages and Coordinator Verbs + +## Status + +Proposed + +## Context + +The Dag processor asks a Lang-SDK runtime two different questions. "Which Dags does this artifact define?" is answered over the messages [ADR-0004](0004-dag-parsing.md) already +defines. "Which task handlers does this artifact register for a `dag_id` Python already owns?" has no answer in those messages, because a `TaskHandler` registration carries no Dag +([ADR-0011](0011-mixed-language-dag-processing.md)). + +This ADR defines the request that carries the second question, the subprocess classes that carry both, and the two parse-side entry points on the coordinator. + +Terms follow the Language SDK spec (`task-sdk/docs/lang-sdk-spec.rst`, spec version `1.0`). + +## Decision + +### One channel shape, two request types + +``` +parse_dag parse_task_handler + + parent process parent process + │ DagFileParseRequest │ TaskHandlerParseRequest + ▼ (ToDagProcessor) ▼ (ToSDKTaskHandlerProcessor) + coordinator (raw byte forward) coordinator (raw byte forward) + ▼ ▼ + runtime runtime + │ DagFileParsingResult │ TaskHandlerParsingResult + ▼ (ToManager) ▼ (ToManager) + parent process parent process +``` + +Both verbs are byte forwarders: the coordinator spawns the runtime, wires `fd 0` to the comm socket, and never decodes the payload. The process that spawned the parse decodes the +reply. They stay two methods, not one `parse(request)`, because a coordinator can serve handlers without serving native Dag parsing. Appendix A has the longer argument. + +### The reply travels on `ToManager` + +``` +_ParseSideResponses = shared tail — same members, same `type` discriminator + ConnectionResult | VariableResult | VariableKeysResult | TaskStatesResult + | PreviousDagRunResult | PreviousTIResult | PrevSuccessfulDagRunResult + | ErrorResponse | OKResponse | XComCountResponse | XComResult + | XComSequenceIndexResult | XComSequenceSliceResult + +ToDagProcessor = DagFileParseRequest | _ParseSideResponses parent → child +ToSDKTaskHandlerProcessor = TaskHandlerParseRequest | _ParseSideResponses parent → child (new) + +ToManager = DagFileParsingResult | TaskHandlerParsingResult child → parent + | GetConnection | GetVariable | … | MaskSecret +``` + +`ToSDKTaskHandlerProcessor` is the only new union; `ToManager` gains one member. The two parent → child unions differ in exactly one member, because the child's questions about +connections, variables and XComs do not depend on which parse it was asked for. + +`ToManager` is named for the process that usually holds the other end, but the role it describes is "whoever spawned this parse". The Dag processor manager fills it for +`DagFileProcessorProcess`; a Dag-parsing child fills it for the two processes below, relaying anything that is not a parsing result up its own `ToManager` channel unchanged. That +relay is only type-safe because both hops speak the same pair, which is the reason not to mint a separate `ToCoordinator`. + +### Message shapes + +``` +class TaskHandlerParseRequest: + file: str # the artifact resolved for this coordinator + dag_ids: list[str] # every Dag in the parsed file with stub tasks that resolved here + bundle_path: Path + bundle_name: str + type: Literal["TaskHandlerParseRequest"] + +class TaskHandlerParsingResult: + fileloc: str + task_handlers: dict[str, list[TaskHandlerDeclaration]] # dag_id → its declarations + import_errors: dict[str, str] | None = None + warnings: list | None = None + type: Literal["TaskHandlerParsingResult"] + +class TaskHandlerDeclaration: + task_id: str + params: list[TaskHandlerParam] # ordered — arg_bindings are positional + +class TaskHandlerParam: + name: str + value_schema: ArgValueSchema | None = None + required: bool # the handler declares no default Review Comment: Where is this used? The flag is not mentioned anywhere else (both ADRs in this PR and previously merged) -- 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]
