Skip to content

Probe a Lang-SDK artifact for its task handlers - #73974

Draft
jason810496 wants to merge 11 commits into
jason/core-taskhandler-refactor/04-task-handler-parse-messagesfrom
jason/core-taskhandler-refactor/05-runtime-parse-transport
Draft

jason810496 wants to merge 11 commits into
jason/core-taskhandler-refactor/04-task-handler-parse-messagesfrom
jason/core-taskhandler-refactor/05-runtime-parse-transport

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

To check stub tasks, the Dag processor has to find the artifact a worker would run, start its runtime inside the Dag-parsing child, and ask it with #73973's messages. This PR adds that probe and the two coordinator hooks it needs. Nothing calls it yet: #73975, #74317 and #73976 implement the hooks for Go, Java and TypeScript, and #74135 calls the probe.

class MyCoordinator(SubprocessCoordinator):
    def _find_task_handler_artifact(self, *, bundle_path, dag_id):
        # The artifact a stub task of dag_id runs, found the way execute_task finds it.
        return ResolvedBundle(bundle_path / "handlers", "2026-10-30")

    def _build_parse_task_handler_command(self, *, path):
        # The command without --comm / --logs, and the runtime's supervisor schema version.
        return [os.fspath(path)], "2026-10-30"


result = LangSDKTaskHandlerProcessorProcess.run(
    coordinator="my-sdk",  # key in [sdk] coordinators
    path=bundle_path / "handlers",
    bundle_path=bundle_path,
    bundle_name="task-handlers",
    artifact_rel_path="handlers",
    logger=log,
    deadline=time.monotonic() + 30,  # optional
)  # every handler the artifact registers; a failure is in result.import_errors
  • The hooks live on SubprocessCoordinator only, and BaseCoordinator is unchanged. A coordinator without them gives an import error that says so.
  • Only a runtime that reports supervisor schema 2026-10-30 or later is probed. That is the in-progress version, so no new date is invented.
  • LangSDKTaskHandlerProcessorProcess is a standalone process class beside Parse native Lang-SDK Dags in the Dag processor with their runtime #74035's LangSDKDagFileProcessorProcess, with no API client. A failed start, a missing or invalid result, or a timeout is an import error keyed by the artifact's path in its bundle.
  • Requests from the runtime that need a client are relayed up the Dag-parsing child's own supervisor channel.
  • The runtime's channel accepts only ToManager, where a DagFileParsingResult is an unhandled request, so a Dag registration is never read as a task handler.
  • The runtime dies with the Dag-parsing child: it runs with PR_SET_PDEATHSIG on Linux, and what it leaves in its process group is killed on exit or timeout.
  • The methods shared with Parse native Lang-SDK Dags in the Dag processor with their runtime #74035's class keep its names and bodies. Lang-SDK runtimes: one transport for native Dag parsing and the task-handler probe #74058 moves both classes onto one transport.
  • A child of a supervised child no longer closes its own new stdin, stdout and stderr when it replaces the handles it inherited. This affects every supervised child, and the probe runs in one.
  • ADR-0012 and ADR-0013 describe the probe as built.

Known limits

  • Where the Dag processor starts children with exec instead of fork (the macOS default), the probe inherits the Dag-parsing child's ORM-blocking environment and dies at import airflow. The check then logs a warning and leaves the stub tasks unchecked. Linux is not affected.
  • The parent death signal reaches only the runtime itself. If the Dag-parsing child is killed abruptly, processes the runtime started survive, so a runtime must exec and leave no children (Lang-SDK task handlers: smaller follow-ups #74059).

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: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/05 runtime parse transport Probe a Lang-SDK artifact for its task handlers Sep 30, 2026
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/05-runtime-parse-transport branch from ff05737 to 99e8fa5 Compare October 1, 2026 06:26
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/05-runtime-parse-transport branch 2 times, most recently from 11b1aa6 to 0b77a65 Compare October 1, 2026 10:55
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/05-runtime-parse-transport branch from 0b77a65 to f6918d8 Compare October 1, 2026 14:44
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/05-runtime-parse-transport branch from f6918d8 to 7a8bab8 Compare October 2, 2026 00:54
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/05-runtime-parse-transport branch from 7a8bab8 to 69e2b80 Compare October 2, 2026 05:48
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/05-runtime-parse-transport branch from 69e2b80 to b45c6ec Compare October 2, 2026 11:36
A supervised child reopens stdin, stdout and stderr on its sockets and
replaces sys.stdin and friends. When that child supervises a child of its
own, as a Dag-parsing child starting a task-handler probe will, the
grandchild replaces handles that already wrap fds 0-2, and dropping them
closed the fds it had just dup2'd, so the grandchild lost its streams.
The handles no longer own their fds; the interpreter's original ones are
kept alive by sys.__stdin__ and friends and were never closed here.
The Dag processor will ask a coordinator's runtime which task handlers an artifact registers (ADR-0011, ADR-0012). SubprocessCoordinator.parse_task_handler replaces the calling process with that runtime, which _build_parse_task_handler_command builds, taking parse_dag's steps and helpers (#74042) in the same order.

On Linux the runtime also gets PR_SET_PDEATHSIG, so it dies with the process that reads its sockets, which the manager kills by pid when a Dag-parsing child times out. BaseCoordinator gains nothing: only a subprocess coordinator can answer.
The parent death signal is only sent for a parent that exits after it is set. A parent that exited between the schema-version report and the prctl call left the runtime with no signal to come and nobody to read its sockets. The child now checks that its parent is unchanged once the signal is set.
The Dag processor will probe the artifact that a stub task actually runs, so a check sees what the worker will run. _find_task_handler_artifact finds it for a Dag the way execute_task does. The coordinators of each SDK implement it in later changes.
A runtime built before task handler parsing cannot answer the request, so starting it only produces a failure. TASK_HANDLER_PARSING_SCHEMA_VERSION is the in-progress supervisor schema version, as native's Java Dag-parsing gate uses, because contributors never add a version date. parse_task_handler refuses an older runtime, and the Dag processor will skip such artifacts before probing.
The Dag-file parse will ask a coordinator's runtime which task handlers an artifact registers (ADR-0011, ADR-0012). LangSDKTaskHandlerProcessorProcess sends TaskHandlerParseRequest and collects one TaskHandlerParsingResult. A failed start, a missing result, an invalid frame or message, or a timeout is an import error on it. Nothing calls it yet.

It is written like LangSDKDagFileProcessorProcess (#74035) and reuses its schema-version message and import timeout, so #74058 can move both onto one transport. It runs in a Dag-parsing child, which has no API client, so the runtime's requests are relayed up that child's supervisor channel. The runtime's channel takes ToManager messages only, and only this class handles TaskHandlerParsingResult.
The frame reader reads a frame's payload in the same callback as its length header. On the blocking comm socket, a runtime that stopped after the header blocked the loop in that read, so the import timeout never fired and the runtime was never killed. Reads now return at once and the reader completes the frame on a later callback; replies still go out on the blocking socket.
A process the runtime started can keep its output open after the runtime exits, so the probe does not finish and that process keeps running. As in LangSDKDagFileProcessorProcess (#74035), the grace kill after the result, the timeout kill and close() now kill the runtime's process group once it has exited, unless its pid was reused.
Only the fork itself was guarded. A failure while registering the listeners or sending the start request left the child waiting for that request and both listeners open, and the standalone run does not guard start() either.
The Dag-file parse will probe its artifacts one after another, and the manager kills the parse child at dag_file_processor_timeout. A probe bounded only by its import timeout could outlast the child, and the parse would send no result at all. run() now also stops at a deadline the caller passes: a probe still running then is killed, and its result is an import error unless the runtime had already answered.
… ADR-0013

ADR-0012 still named BaseParsingProcess, SDKTaskHandlerProcessorProcess and its entry point, put both parse verbs on BaseCoordinator, ran both Lang-SDK processes in a Dag-parsing child, and described the coordinator as forwarding bytes. It now names the classes the code defines, puts the verbs and hooks on SubprocessCoordinator, and says the runtime connects back to the process that asked. ADR-0013's probe step uses the new names.
@jason810496
jason810496 force-pushed the jason/core-taskhandler-refactor/05-runtime-parse-transport branch from b45c6ec to 4a12877 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/05-runtime-parse-transport branch from 4a12877 to 6129fc1 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

area:coordinator Coordinator: The interface to spawn Lang-SDK subprocesses area:DAG-processing area:task-sdk

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant