Skip to content

Go SDK: answer TaskHandlerParseRequest with registered handlers - #73975

Draft
jason810496 wants to merge 6 commits into
jason/core-taskhandler-refactor/05-runtime-parse-transportfrom
jason/core-taskhandler-refactor/06-go-sdk-task-handler-parse
Draft

jason810496 wants to merge 6 commits into
jason/core-taskhandler-refactor/05-runtime-parse-transportfrom
jason/core-taskhandler-refactor/06-go-sdk-task-handler-parse

Conversation

@jason810496

@jason810496 jason810496 commented Sep 30, 2026 •

Copy link
Copy Markdown
Member

Stack (bottom to top), on native #74037 (stack #74170): #73973, #73974, #73975, #74317, #73976, #74067, #74135, #74140, #74141, #73971, #74030, #74031, #74032, #74136, #74137, #74138, #74139

A Go bundle now answers #73974's probe with every handler registered through airflow.TaskHandler and how its params bind, and ExecutableCoordinator can find and start the binary a stub task would run. Until now a Go bundle exited on a TaskHandlerParseRequest as an unknown first frame. Nothing calls the probe yet: #74135 does.

type LoadInput struct {
    Table string `arg:"table_name"`
    Limit int
}

bundle.Register(
    airflow.TaskHandler("etl", "extract", func(actx airflow.Context, region string, day *time.Time) error { ... }),
    airflow.TaskHandler("etl", "load", func(actx airflow.Context, in LoadInput) error { ... }),
)

answers:

TaskHandlerParsingResult(
    fileloc="/bundles/go/etl",
    task_handlers={"etl": [
        TaskHandlerDeclaration(task_id="extract", binding="positional", params=[
            TaskHandlerParam(name=None, value_schema={"type": "string"}),
            TaskHandlerParam(name=None, value_schema={
                "anyOf": [{"type": "string", "format": "date-time"}, {"type": "null"}]}),
        ]),
        TaskHandlerDeclaration(task_id="load", binding="named", params=[
            TaskHandlerParam(name="table_name", exact_name=True, value_schema={"type": "string"}),
            TaskHandlerParam(name="Limit", value_schema={"type": "integer", "format": "int64"}),
        ]),
    ]},
)
  • Only airflow.TaskHandler registrations are declared. Tasks of a Dag built with airflow.Dag are never task handlers, so the answer leaves them out.
  • Binding follows how Go decodes arguments. Flat params are positional and nameless, because Go reflection has no parameter names. A lone struct is named: arg: tags match exactly, other fields ignoring case and _. A lone struct with no bindable fields, such as time.Time, is named with params=None, so only the handler's presence is checked.
  • value_schema states only what decoding into the Go type enforces, in the vocabulary build_arg_bindings emits for Python annotations. Pointers, slices and maps accept null. any, custom json.Unmarshaler types, json.Number and []byte get no schema, because they accept more than one JSON shape.
  • Serve waits for the parent to acknowledge the answer, so the parent has the result before the process exits.
  • The probe checks the binary a worker would run: _find_task_handler_artifact is the worker's own scan, and _build_parse_task_handler_command runs the same trailer and hash check a task gets and returns an absolute path, so exec never searches PATH.
  • The executable bundle scan now walks the Dag bundle in sorted path order, as the TypeScript scan does, so the Dag processor and every worker pick the same binary. A worker's pick changes only when two binaries list the same Dag: the first usable one in sorted order now wins.
  • The language SDK spec gains the task handler parse lifecycle, and its version stays 1.0.

Compatibility: a bundle packed from main's Go SDK before this change reports 2026-10-30 but cannot answer, so its probe fails and #74135 logs a warning. Repacking fixes it. Bundles from the released betas report 2026-06-16 and are not probed.


Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)

Generated-by: Claude Code (Opus 5.5) following the guidelines

@boring-cyborg boring-cyborg Bot added area:coordinator Coordinator: The interface to spawn Lang-SDK subprocesses area:DAG-processing area:dev-tools area:go-sdk area:task-sdk labels Sep 30, 2026
@jason810496
jason810496 added this pull request to stack #73978 September 30, 2026 17:59
@jason810496 jason810496 changed the title jason/core taskhandler refactor/06 go sdk task handler parse Go SDK: answer TaskHandlerParseRequest with registered handlers Sep 30, 2026
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/06-go-sdk-task-handler-parse branch from 337387a to b410623 Compare October 1, 2026 06:26
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/06-go-sdk-task-handler-parse branch from b410623 to f17015e Compare October 1, 2026 06:30
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/06-go-sdk-task-handler-parse branch from f17015e to 421cb23 Compare October 1, 2026 10:55
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/06-go-sdk-task-handler-parse branch 2 times, most recently from 78995cf to b1151dc Compare October 2, 2026 00:54
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/06-go-sdk-task-handler-parse branch from b1151dc to 0fb5c4c Compare October 2, 2026 05:48
The Dag processor will check each stub task against the Go handler it runs, so the binding plan has to say how the handler's params bind and which values they accept. Flat params are positional and nameless. A lone struct, tagged or not, is named, and arg tags match exactly. Each param gets a JSON Schema in the vocabulary build_arg_bindings emits, stating only what decoding enforces. Params is an empty list rather than null when there are none, because null means the params cannot be listed.
The Dag processor asks the binary a stub task would run which task handlers it registers, and a Go bundle exits on that first frame as an unknown message. Serve now replies with one TaskHandlerParsingResult listing every handler registered with airflow.TaskHandler, keyed by Dag id in registration order, so the answer depends only on the bundle. It waits for the parent's acknowledgement before returning, so the parent has the result before the process exits. A Dag from airflow.Dag is never a task handler and is left out. Empty params and task_handlers go out as empty values, because the parent rejects a null task_handlers and reads null params as "cannot list".
The Dag processor will probe the binary that the worker's scan picks, and that scan followed each directory's filesystem order. When two binaries list the same Dag, two hosts could pick different ones. The scan now uses walk_files (#73126), as the TypeScript scan does, so the first match in sorted path order wins on every host. A worker runs another binary than before only when two binaries list the same Dag.
The probe needs each coordinator to find the artifact a stub task runs and to start its runtime. ExecutableCoordinator finds the bundle with the scan that execute_task uses, so the Dag processor checks the binary a worker would run, and returns its resolved path, as the task command does. The parse command runs the trailer and binary hash check a task's bundle gets and returns the absolute path, so exec never searches PATH for a relative name.
The probe tests fake the runtime, so nothing showed that a real Go bundle answers the way the Dag processor parses it. The test packs go-sdk/example/bundle, finds it as a stub task of simple_dag would, and probes it. It needs a Go toolchain, so it runs only with AIRFLOW_LANG_SDK_REAL_PROBE_TESTS=1.
The Go SDK is the first runtime to answer TaskHandlerParseRequest, and the Java and TypeScript runtimes follow the spec and the contributor guide. Neither said which coordinator hooks find and start the runtime, that the first frame can be a parse request, what the answer must contain, or what the Dag processor does with a mismatch under each binding. The spec stays at 1.0: the lifecycle adds steps and renames no term.
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/06-go-sdk-task-handler-parse branch from ccecc81 to cf557b9 Compare October 5, 2026 18:12
@jason810496
jason810496 removed this pull request from stack #73978 October 6, 2026 02:27
@jason810496
jason810496 added this pull request to stack #74318 October 6, 2026 02:35
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/06-go-sdk-task-handler-parse branch from cf557b9 to 04698db Compare October 6, 2026 02:44

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant