This is an automated email from the ASF dual-hosted git repository.

jason810496 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new 8cdb183d09e TS SDK: embed the entrypoint source in the Dag bundle 
(#73127)
8cdb183d09e is described below

commit 8cdb183d09e4fc27ff0fef70f839d446022b61fe
Author: Jason(Zhe-You) Liu <[email protected]>
AuthorDate: Sat Sep 19 11:25:51 2026 +0800

    TS SDK: embed the entrypoint source in the Dag bundle (#73127)
    
    A packed bundle now ships the entry module as its author wrote it, in a
    `/*# airflowSource ... #*/` block comment with its own byte range and 
SHA-256
    digest in the layout header. The shipped code region is minified, so it is 
not
    the code anyone wrote and there was nothing readable for the Airflow UI to 
show
    for a natively authored TypeScript Dag. Only the entry module is embedded, 
not
    the module graph behind it; ADR-0006 declined multi-file source display for
    mixed-language Dags.
    
    A block comment, not the line comments the layout and metadata use, because 
the
    entrypoint spans the lines it was written on. That makes line terminators 
legal
    inside the region and leaves `*/` as the one sequence that must not appear: 
it
    would end the comment where Node reads the file, putting the rest of the 
payload
    into executable position while the layout still calls those bytes source and
    both digests still match. The packer inserts a `\` between the two 
characters,
    and escapes `*\` the same way so the transformation is reversible; the 
reader
    reverses it and rejects a bundle carrying the sequence unescaped rather than
    vouch for one.
    
    The cross-language golden fixture now carries regions shaped like what the
    packer ships, instead of the raw contents of a `.ts` fixture: a code region
    nobody could read, and a source region that exercises both escape branches. 
The
    escape scheme is a byte-level agreement between the TypeScript encoder and 
the
    Python reader, and only an artifact the encoder actually produced can 
establish
    it; the Python-side helper that escapes source for other tests is a
    reimplementation and cannot. The fixture is also excluded from Prettier, 
which
    would otherwise reformat it and invalidate its digests.
    
    The metadata `dags` mapping is renamed to `task_handlers`, which is what a
    TypeScript bundle actually provides: handlers for Dags declared elsewhere, 
not
    Dag definitions of its own. The mapping stays keyed by Dag ID.
    
    airflow_bundle_metadata_version stays at 1.0, redefined in place because the
    TypeScript packing workflow is unreleased. A source-less bundle still fails
    closed, on the now-required `source` layout section, before any version 
check
    is reached.
---
 .rat-excludes                                      |   3 +
 .../language-sdks/typescript.rst                   |   4 +
 docs/spelling_wordlist.txt                         |   1 +
 task-sdk/docs/ts-bundle-spec.rst                   |  79 ++++++++---
 .../sdk/coordinators/node/_bundle_reader.py        |  99 +++++++++++--
 .../coordinators/node/_bundle_test_utils.py        |  23 ++-
 .../coordinators/node/test_bundle_reader.py        | 143 +++++++++++++++++--
 ts-sdk/.prettierignore                             |   3 +
 ts-sdk/README.md                                   |   8 +-
 ts-sdk/src/cli/bundle-encoder.ts                   |  72 ++++++++--
 ts-sdk/src/cli/pack.ts                             |  25 ++--
 ts-sdk/src/cli/validate.ts                         |   6 +-
 ts-sdk/src/coordinator/manifest.ts                 |  17 ++-
 ts-sdk/tests/cli/fixtures/bundle-v1.min.mjs        |  18 +++
 ts-sdk/tests/cli/fixtures/bundle-v1.mjs            |  24 ----
 ts-sdk/tests/cli/pack.test.ts                      | 155 +++++++++++++++++----
 ts-sdk/tests/cli/validate.test.ts                  |   4 +-
 ts-sdk/tests/coordinator/runtime-manifest.test.ts  |  14 +-
 18 files changed, 550 insertions(+), 148 deletions(-)

diff --git a/.rat-excludes b/.rat-excludes
index 4dd747ada12..dc4e0ffcd6a 100644
--- a/.rat-excludes
+++ b/.rat-excludes
@@ -328,6 +328,9 @@ www-hash.txt
 # Vendored-in code
 /src/airflow/providers/google/_vendor/*
 
+# TypeScript bundle golden fixture: its bytes are digest-pinned, so no license 
header can be added
+/ts-sdk/tests/cli/fixtures/bundle-v1.min.mjs
+
 # Java SDK build outputs
 /java-sdk/bin/*
 /java-sdk/build/*
diff --git 
a/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst 
b/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst
index d374225fca4..ad033147bfe 100644
--- a/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst
+++ b/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst
@@ -326,6 +326,10 @@ read or edit in place. The ``/*! */`` license banners of 
bundled dependencies ar
 a function name, so minified names are safe: a Dag and a task are named by the 
string ids their registration
 states, and a handler is dispatched by reference.
 
+Because the shipped code is not the code anyone wrote, the packer also embeds 
the entry module verbatim in a
+``/*# airflowSource ... #*/`` block comment, verified by its own digest, so 
Airflow has something readable to
+display for the Dag. Only the entry module is embedded, not the modules it 
imports.
+
 ``esbuild`` is an optional peer dependency: packing is build-time only, so the 
runtime install of
 ``apache-airflow-ts-sdk`` skips it, and it must be installed separately before 
running ``airflow-ts-pack``.
 
diff --git a/docs/spelling_wordlist.txt b/docs/spelling_wordlist.txt
index f35a04c3b16..dedf87bc4fe 100644
--- a/docs/spelling_wordlist.txt
+++ b/docs/spelling_wordlist.txt
@@ -1843,6 +1843,7 @@ uncompress
 undeploy
 undeployed
 Undeploys
+unescaped
 unicode
 unicodecsv
 unindent
diff --git a/task-sdk/docs/ts-bundle-spec.rst b/task-sdk/docs/ts-bundle-spec.rst
index 82147f95e5c..705227efa55 100644
--- a/task-sdk/docs/ts-bundle-spec.rst
+++ b/task-sdk/docs/ts-bundle-spec.rst
@@ -32,16 +32,16 @@ that suffix.
 Container
 ---------
 
-The bundle remains an ECMAScript module that runs directly with ``node 
bundle.min.mjs``. It has three regions:
+The bundle remains an ECMAScript module that runs directly with ``node 
bundle.min.mjs``. It has four regions:
 
 .. code-block:: text
 
     //# airflowBundle=<compact JSON layout>\n
     //# airflowMetadata=<compact JSON>\n
+    /*# airflowSource\n<escaped entrypoint source>\n#*/\n
     <minified, bundled ECMAScript code>
 
-The layout comes first so readers can locate and verify the other regions. The
-current format has no embedded source region.
+The layout comes first so readers can locate and verify the other regions.
 
 Layout Header
 -------------
@@ -60,6 +60,11 @@ The ``airflowBundle`` payload is a compact UTF-8 JSON object:
         "start": "0000000000000300",
         "end": "0000000000000400",
         "sha256": 
"123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef0"
+      },
+      "source": {
+        "start": "0000000000000412",
+        "end": "00000000000003f0",
+        "sha256": 
"23456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef01"
       }
     }
 
@@ -75,13 +80,41 @@ including this physical framing and the decoded metadata 
schema. A reader parses
 can locate and verify that version.
 
 The metadata range points to the UTF-8 JSON payload only, excluding the 
JavaScript comment marker and newline. Its
-digest therefore covers the exact JSON bytes stored in that range. The code 
range covers every byte after the
-metadata line through the end of the file, and its digest covers those raw 
JavaScript bytes.
-
-The file begins with the layout line. The metadata marker immediately follows 
that line, and exactly one newline
-separates the metadata payload from the code range. These prescribed framing 
bytes are outside the hashed metadata
-and code ranges, and no additional bytes are permitted before, between, or 
after them. Post-pack formatters,
-compressors, source-map injectors, and other tools that rewrite the bundle 
invalidate the offsets or digests.
+digest therefore covers the exact JSON bytes stored in that range. The source 
range likewise points to the escaped
+entrypoint payload only, excluding the ``/*# airflowSource\n`` opener and the 
``\n#*/\n`` closer. The code range
+covers every byte after the source comment through the end of the file, and 
its digest covers those raw JavaScript
+bytes.
+
+The file begins with the layout line. The metadata marker immediately follows 
that line, exactly one newline
+separates the metadata payload from the source comment's opener, and the 
source comment's closer is immediately
+followed by the code range. These prescribed framing bytes are outside the 
hashed ranges, and no additional bytes are
+permitted before, between, or after them. Post-pack formatters, compressors, 
source-map injectors, and other tools
+that rewrite the bundle invalidate the offsets or digests.
+
+Unlike the metadata range, the source range's length is declared rather than 
derivable from a newline, so a reader
+pins it by checking that the prescribed opener and closer sit exactly where 
the declared range implies.
+
+Source
+------
+
+The source region carries the entrypoint as its author wrote it, so the 
Airflow UI has something readable to show
+for a natively authored TypeScript Dag. The shipped code region is minified 
and is not the code anyone wrote. Only
+the entrypoint is embedded, not the module graph behind it:
+`ADR-0006 
<https://github.com/apache/airflow/blob/main/airflow-core/adr/lang-sdk/0006-no-lang-sdk-source-display.md>`__
+declined multi-file source display for mixed-language Dags.
+
+It is a block comment rather than the line comments the layout and metadata 
use, because the entrypoint spans the
+lines it was written on and a ``//`` comment would end at the first of them. 
Line terminators, including ``\r`` and
+U+2028/U+2029, are therefore legal inside the region and need no escaping.
+
+``*/`` is the one sequence that must not appear: it would end the comment 
where Node reads the file, putting the rest
+of the payload into executable position while the layout still calls those 
bytes source and both digests still match.
+The packer escapes it by inserting a ``\`` between the two characters, and 
escapes ``*\`` the same way so the
+transformation is reversible. A reader recovers the entrypoint by dropping the 
``\`` that follows a ``*`` and keeping
+the character behind it, and MUST reject a bundle whose source range contains 
an unescaped ``*/``.
+
+The region is capped at 1 MiB. Larger entrypoints should move code into 
imported modules, which are bundled into the
+code region as usual.
 
 Metadata
 --------
@@ -98,7 +131,7 @@ The ``airflowMetadata`` payload is compact UTF-8 JSON with 
this logical shape:
         "supervisor_schema_version": "2026-06-16"
       },
       "source": "main.ts",
-      "dags": {
+      "task_handlers": {
         "example": {
           "tasks": ["extract", "load"]
         }
@@ -109,8 +142,13 @@ The packer serializes this object without insignificant 
whitespace and escapes t
 separators (U+2028 and U+2029), keeping it in one newline-terminated 
JavaScript comment without a second encoding
 layer. The SHA-256 digest detects changes to the exact serialized bytes.
 
-The coordinator uses the ``dags`` keys to choose a bundle for a task instance. 
The ``source`` value is a logical
-authoring name only, not embedded source content, and it is not used to 
execute the bundle.
+``task_handlers`` is keyed by Dag ID and lists the task IDs the bundle handles 
for each. It is named for what a
+TypeScript bundle actually provides: handlers for Dags declared elsewhere, not 
Dag definitions of its own. The
+coordinator uses its keys to choose a bundle for a task instance.
+
+The metadata ``source`` value is the logical authoring name displayed for the 
Dag, a filename rather than content.
+The source region named in the layout header is what carries the content. The 
two are separate fields in separate
+documents, and neither is used to execute the bundle.
 
 Reader and Selection Algorithm
 ------------------------------
@@ -122,14 +160,15 @@ For each candidate in ``bundles_root``, the coordinator:
    a filesystem returns entries in. Directories are deduplicated by ``(st_dev, 
st_ino)``, so a symlink loop
    terminates the walk instead of exhausting the interpreter stack.
 2. Reads a bounded first line and decodes the named metadata and code ranges.
-3. Reads the bounded metadata line and checks that the declared metadata and 
code ranges exactly match their
-   physical locations and the file size.
-4. Computes SHA-256 for both ranges before parsing or using metadata.
-5. Confirms with ``fstat`` that the open file did not change during 
verification.
-6. Parses metadata and requires a supported bundle contract major version from
+3. Reads the bounded metadata line and checks that the declared metadata range 
matches its physical location.
+4. Reads the bounded source region, checks that its prescribed opener and 
closer frame the declared range, rejects an
+   unescaped ``*/`` inside it, and checks that the declared code range matches 
its physical location and file size.
+5. Computes SHA-256 for all three ranges before parsing or using metadata.
+6. Confirms with ``fstat`` that the open file did not change during 
verification.
+7. Parses metadata and requires a supported bundle contract major version from
    ``airflow_bundle_metadata_version``.
-7. Skips the verified bundle if its ``dags`` mapping does not contain the 
requested ``dag_id``.
-8. Resolves the supervisor schema version and selects the first usable match.
+8. Skips the verified bundle if its ``task_handlers`` mapping does not contain 
the requested ``dag_id``.
+9. Resolves the supervisor schema version and selects the first usable match.
 
 A missing, unrelated, unreadable, malformed, corrupt, or incompatible earlier 
candidate does not prevent selection
 of a later usable match. When more than one usable bundle declares the same 
Dag, the first configured match wins. If
diff --git a/task-sdk/src/airflow/sdk/coordinators/node/_bundle_reader.py 
b/task-sdk/src/airflow/sdk/coordinators/node/_bundle_reader.py
index a03576aeac7..db766d1433a 100644
--- a/task-sdk/src/airflow/sdk/coordinators/node/_bundle_reader.py
+++ b/task-sdk/src/airflow/sdk/coordinators/node/_bundle_reader.py
@@ -18,8 +18,8 @@
 """
 Read and verify TypeScript Dag bundles.
 
-File order: layout comment, metadata comment, then executable JavaScript.
-Read the headers, verify the section digests, and return coordinator metadata.
+File order: layout comment, metadata comment, source block comment, then 
executable JavaScript.
+Read the headers, verify the section digests, and return coordinator metadata 
or the source.
 """
 
 from __future__ import annotations
@@ -44,7 +44,13 @@ _LAYOUT_COMMENT_PREFIX = b"//# airflowBundle="
 _MAX_LAYOUT_LINE_BYTES = 4096
 _METADATA_COMMENT_PREFIX = b"//# airflowMetadata="
 _MAX_METADATA_LINE_BYTES = 1024 * 1024
+# A block comment, because the entrypoint spans more than one line.
+_SOURCE_COMMENT_OPEN = b"/*# airflowSource\n"
+_SOURCE_COMMENT_CLOSE = b"\n#*/\n"
+_MAX_SOURCE_BYTES = 1024 * 1024
 _SUPPORTED_BUNDLE_MAJOR_VERSION = 1
+# Reverses the packer's escaping of ``*/``, the only sequence that could close 
the comment early.
+_SOURCE_ESCAPE = re.compile(rb"\*\\([\\/])")
 
 # Bound hashing memory and process-local cache growth independently of the 
format.
 _HASH_CHUNK_BYTES = 1024 * 1024
@@ -63,9 +69,10 @@ class _DeclaredSection:
 
 @attrs.define(frozen=True)
 class _BundleLayout:
-    """The metadata and executable sections described by the layout comment."""
+    """The metadata, source, and executable sections described by the layout 
comment."""
 
     metadata: _DeclaredSection
+    source: _DeclaredSection
     code: _DeclaredSection
 
 
@@ -75,6 +82,7 @@ class _DigestCacheKey:
 
     path: str
     metadata: _DeclaredSection
+    source: _DeclaredSection
     code: _DeclaredSection
     device: int
     inode: int
@@ -88,6 +96,7 @@ class _ComputedDigests:
     """SHA-256 digests calculated from the actual section bytes."""
 
     metadata: bytes
+    source: bytes
     code: bytes
 
 
@@ -106,6 +115,25 @@ class BundleMetadata:
 
 def read_bundle(bundle_path: pathlib.Path) -> BundleMetadata:
     """Read and verify one exact TypeScript bundle file."""
+    # Interpret metadata only after checking its serialized bytes against the 
declared digest.
+    return 
_parse_bundle_metadata(_read_verified_payloads(bundle_path).metadata)
+
+
+def read_bundle_source(bundle_path: pathlib.Path) -> str:
+    """Return the entrypoint source ``airflow-ts-pack`` embedded in 
*bundle_path*."""
+    return _decode_source(_read_verified_payloads(bundle_path).source)
+
+
[email protected](frozen=True)
+class _VerifiedPayloads:
+    """The metadata and source payloads of a bundle whose digests have been 
checked."""
+
+    metadata: bytes
+    source: bytes
+
+
+def _read_verified_payloads(bundle_path: pathlib.Path) -> _VerifiedPayloads:
+    """Open *bundle_path* once, check its framing and digests, and return its 
payloads."""
     try:
         bundle_file = bundle_path.open("rb")
     except OSError as exc:
@@ -118,13 +146,19 @@ def read_bundle(bundle_path: pathlib.Path) -> 
BundleMetadata:
         except OSError as exc:
             raise OSError(f"cannot read {bundle_path.name}: {exc}") from exc
 
-        layout, metadata_payload = _read_bundle_headers(
+        layout, payloads = _read_bundle_headers(
             bundle_file, path=bundle_path, file_size=initial_file_info.st_size
         )
         _verify_integrity(bundle_file, path=bundle_path, layout=layout, 
initial_file_info=initial_file_info)
 
-    # Interpret metadata only after checking its serialized bytes against the 
declared digest.
-    return _parse_bundle_metadata(metadata_payload)
+    return payloads
+
+
+def _decode_source(payload: bytes) -> str:
+    try:
+        return _SOURCE_ESCAPE.sub(rb"*\1", payload).decode("utf-8")
+    except UnicodeDecodeError as exc:
+        raise ValueError(f"embedded airflow source is not valid UTF-8: {exc}") 
from exc
 
 
 class _BundleDigestCache:
@@ -191,6 +225,7 @@ def _parse_layout(payload: bytes) -> _BundleLayout:
         raise ValueError("embedded airflow bundle layout must contain a 
mapping")
     return _BundleLayout(
         metadata=_parse_section(layout, "metadata"),
+        source=_parse_section(layout, "source"),
         code=_parse_section(layout, "code"),
     )
 
@@ -266,6 +301,9 @@ def _compute_stable_digests(
         metadata=_hash_region(
             bundle_file, start=layout.metadata.start, end=layout.metadata.end, 
path=path, section="metadata"
         ),
+        source=_hash_region(
+            bundle_file, start=layout.source.start, end=layout.source.end, 
path=path, section="source"
+        ),
         code=_hash_region(
             bundle_file, start=layout.code.start, end=layout.code.end, 
path=path, section="code"
         ),
@@ -290,6 +328,7 @@ def _verify_integrity(
     cache_key = _DigestCacheKey(
         path=os.fspath(path),
         metadata=layout.metadata,
+        source=layout.source,
         code=layout.code,
         device=initial_file_info.st_dev,
         inode=initial_file_info.st_ino,
@@ -307,6 +346,7 @@ def _verify_integrity(
 
     for section, computed_digest, declared_digest in (
         ("metadata", computed_digests.metadata, layout.metadata.sha256),
+        ("source", computed_digests.source, layout.source.sha256),
         ("code", computed_digests.code, layout.code.sha256),
     ):
         if computed_digest != declared_digest:
@@ -335,19 +375,20 @@ def _parse_bundle_metadata(payload: bytes) -> 
BundleMetadata:
             f"unsupported airflow bundle metadata version {value!r}; "
             f"this runtime supports major version 
{_SUPPORTED_BUNDLE_MAJOR_VERSION}"
         )
-    dags = metadata.get("dags")
-    if not isinstance(dags, dict):
-        raise ValueError("embedded airflow metadata must contain a dags 
mapping")
+    # Keyed by Dag, but named for what the bundle provides: handlers for Dags 
declared elsewhere.
+    task_handlers = metadata.get("task_handlers")
+    if not isinstance(task_handlers, dict):
+        raise ValueError("embedded airflow metadata must contain a 
task_handlers mapping")
     return BundleMetadata(
-        dag_ids=frozenset(dags),
+        dag_ids=frozenset(task_handlers),
         supervisor_schema_version=extract_supervisor_schema_version(metadata),
     )
 
 
 def _read_bundle_headers(
     bundle_file: BinaryIO, *, path: pathlib.Path, file_size: int
-) -> tuple[_BundleLayout, bytes]:
-    """Read both comment payloads and check their declared ranges against the 
file."""
+) -> tuple[_BundleLayout, _VerifiedPayloads]:
+    """Read every comment payload and check their declared ranges against the 
file."""
     layout_payload = _read_prefixed_line(
         bundle_file,
         path=path,
@@ -368,9 +409,39 @@ def _read_bundle_headers(
     layout_line_size = len(_LAYOUT_COMMENT_PREFIX) + len(layout_payload) + 1
     metadata_start = layout_line_size + len(_METADATA_COMMENT_PREFIX)
     metadata_end = metadata_start + len(metadata_payload)
-    code_start = metadata_end + 1
     if (layout.metadata.start, layout.metadata.end) != (metadata_start, 
metadata_end):
         raise ValueError("bundle layout metadata offsets do not match the 
metadata section")
+    source_payload = _read_source_region(
+        bundle_file, path=path, start=metadata_end + 1 + 
len(_SOURCE_COMMENT_OPEN), layout=layout
+    )
+    code_start = layout.source.end + len(_SOURCE_COMMENT_CLOSE)
     if (layout.code.start, layout.code.end) != (code_start, file_size):
         raise ValueError("bundle layout code offsets do not match the 
executable section")
-    return layout, metadata_payload
+    return layout, _VerifiedPayloads(metadata=metadata_payload, 
source=source_payload)
+
+
+def _read_source_region(
+    bundle_file: BinaryIO, *, path: pathlib.Path, start: int, layout: 
_BundleLayout
+) -> bytes:
+    """Read the source block comment and check that it frames the declared 
range exactly."""
+    # The length is declared, not derivable, so the comment markers pin the 
range to real bytes.
+    if layout.source.start != start:
+        raise ValueError("bundle layout source offsets do not match the source 
section")
+    length = layout.source.end - start
+    if length > _MAX_SOURCE_BYTES:
+        raise ValueError(f"embedded airflow source exceeds {_MAX_SOURCE_BYTES} 
bytes")
+    try:
+        opener = bundle_file.read(len(_SOURCE_COMMENT_OPEN))
+        payload = bundle_file.read(length)
+        closer = bundle_file.read(len(_SOURCE_COMMENT_CLOSE))
+    except OSError as exc:
+        raise OSError(f"cannot read {path.name}: {exc}") from exc
+    if opener != _SOURCE_COMMENT_OPEN:
+        raise ValueError(f"{path.name} has no embedded airflow source after 
its metadata")
+    if len(payload) != length or closer != _SOURCE_COMMENT_CLOSE:
+        raise ValueError("embedded airflow source is not closed by its block 
comment")
+    # An unescaped ``*/`` would end the comment where Node reads the file,
+    # so bytes the layout calls source would execute while both digests still 
match.
+    if b"*/" in payload:
+        raise ValueError("embedded airflow source contains an unescaped block 
comment terminator")
+    return payload
diff --git a/task-sdk/tests/task_sdk/coordinators/node/_bundle_test_utils.py 
b/task-sdk/tests/task_sdk/coordinators/node/_bundle_test_utils.py
index 0ec0a13cf7b..ee63b5b2df9 100644
--- a/task-sdk/tests/task_sdk/coordinators/node/_bundle_test_utils.py
+++ b/task-sdk/tests/task_sdk/coordinators/node/_bundle_test_utils.py
@@ -21,11 +21,14 @@ from __future__ import annotations
 import hashlib
 import json
 import pathlib
+import re
 
 SCHEMA_VERSION = "2026-06-16"
 BUNDLE_NAME = "bundle.min.mjs"
 LAYOUT_PREFIX = b"//# airflowBundle="
 METADATA_PREFIX = b"//# airflowMetadata="
+SOURCE_OPEN = b"/*# airflowSource\n"
+SOURCE_CLOSE = b"\n#*/\n"
 OFFSET_WIDTH = 16
 
 
@@ -41,7 +44,7 @@ def metadata_json(
             "supervisor_schema_version": schema_version,
         },
         "source": "main.ts",
-        "dags": {dag_id: {"tasks": ["test_task"]} for dag_id in dag_ids},
+        "task_handlers": {dag_id: {"tasks": ["test_task"]} for dag_id in 
dag_ids},
     }
     if metadata_version is not None:
         metadata = {"airflow_bundle_metadata_version": metadata_version, 
**metadata}
@@ -61,6 +64,11 @@ def _layout_line(layout: dict[str, object]) -> bytes:
     return LAYOUT_PREFIX + payload + b"\n"
 
 
+def escape_source(source: bytes) -> bytes:
+    """Escape a block-comment terminator the way the TypeScript encoder 
does."""
+    return re.sub(rb"\*([\\/])", rb"*\\\1", source)
+
+
 def write_bundle(
     root: pathlib.Path,
     *dag_ids: str,
@@ -68,6 +76,8 @@ def write_bundle(
     schema_version: str = SCHEMA_VERSION,
     metadata_version: str | None = "1.0",
     metadata_payload: bytes | None = None,
+    source: bytes = b"export {};\n",
+    source_payload: bytes | None = None,
     name: str = BUNDLE_NAME,
 ) -> pathlib.Path:
     if metadata_payload is None:
@@ -76,27 +86,34 @@ def write_bundle(
             schema_version=schema_version,
             metadata_version=metadata_version,
         )
+    if source_payload is None:
+        source_payload = escape_source(source)
     metadata_line = METADATA_PREFIX + metadata_payload + b"\n"
+    source_region = SOURCE_OPEN + source_payload + SOURCE_CLOSE
     placeholder = _layout_line(
         {
             "code": _section(0, 0, code),
             "metadata": _section(0, 0, metadata_payload),
+            "source": _section(0, 0, source_payload),
         }
     )
     metadata_start = len(placeholder) + len(METADATA_PREFIX)
     metadata_end = metadata_start + len(metadata_payload)
-    code_start = len(placeholder) + len(metadata_line)
+    source_start = len(placeholder) + len(metadata_line) + len(SOURCE_OPEN)
+    source_end = source_start + len(source_payload)
+    code_start = len(placeholder) + len(metadata_line) + len(source_region)
     layout_line = _layout_line(
         {
             "code": _section(code_start, code_start + len(code), code),
             "metadata": _section(metadata_start, metadata_end, 
metadata_payload),
+            "source": _section(source_start, source_end, source_payload),
         }
     )
     assert len(layout_line) == len(placeholder)
 
     bundle = root / name
     bundle.parent.mkdir(parents=True, exist_ok=True)
-    bundle.write_bytes(layout_line + metadata_line + code)
+    bundle.write_bytes(layout_line + metadata_line + source_region + code)
     return bundle
 
 
diff --git a/task-sdk/tests/task_sdk/coordinators/node/test_bundle_reader.py 
b/task-sdk/tests/task_sdk/coordinators/node/test_bundle_reader.py
index 3866028b658..3f3beaeefd7 100644
--- a/task-sdk/tests/task_sdk/coordinators/node/test_bundle_reader.py
+++ b/task-sdk/tests/task_sdk/coordinators/node/test_bundle_reader.py
@@ -31,6 +31,8 @@ from task_sdk.coordinators.node._bundle_test_utils import (
     METADATA_PREFIX,
     OFFSET_WIDTH,
     SCHEMA_VERSION,
+    SOURCE_CLOSE,
+    SOURCE_OPEN,
     metadata_json as _metadata_json,
     mutate_byte as _mutate_byte,
     read_layout as _read_layout,
@@ -40,11 +42,16 @@ from task_sdk.coordinators.node._bundle_test_utils import (
 )
 
 from airflow.sdk.coordinators.node import _bundle_reader as _reader
-from airflow.sdk.coordinators.node._bundle_reader import _digest_cache, 
_hash_region, read_bundle
+from airflow.sdk.coordinators.node._bundle_reader import (
+    _digest_cache,
+    _hash_region,
+    read_bundle,
+    read_bundle_source,
+)
 
 from tests_common.test_utils.paths import AIRFLOW_ROOT_PATH
 
-TYPESCRIPT_V1_FIXTURE = AIRFLOW_ROOT_PATH / "ts-sdk" / "tests" / "cli" / 
"fixtures" / "bundle-v1.mjs"
+TYPESCRIPT_V1_FIXTURE = AIRFLOW_ROOT_PATH / "ts-sdk" / "tests" / "cli" / 
"fixtures" / "bundle-v1.min.mjs"
 
 
 @pytest.fixture(autouse=True)
@@ -60,6 +67,31 @@ class TestBundleReader:
         assert metadata.dag_ids == frozenset({"test_dag"})
         assert metadata.supervisor_schema_version == SCHEMA_VERSION
 
+    def test_reads_source_embedded_by_typescript_encoder(self):
+        # Real encoder output, whose source region needs both escape branches. 
Recovering it
+        # exactly is the cross-language agreement on the escape scheme, which 
the Python helper
+        # used by the other tests cannot establish on its own.
+        assert read_bundle_source(TYPESCRIPT_V1_FIXTURE) == (
+            "/** Handlers for the test Dag. */\n"
+            'import { Bundle, Dag } from "apache-airflow-ts-sdk";\n'
+            "\n"
+            r"const TERMINATOR = /\*\//;"
+            "\n"
+            'const dag = new Dag("test_dag");\n'
+            'dag.task("test_task", async () => TERMINATOR.source);\n'
+            "\n"
+            "await new Bundle(dag).serve();\n"
+        )
+
+    def test_accepts_block_comment_terminator_inside_the_code_region(self):
+        # Rejecting a terminator in the code region would make every real 
bundle unreadable.
+        fixture = TYPESCRIPT_V1_FIXTURE.read_bytes()
+        layout = _read_layout(TYPESCRIPT_V1_FIXTURE)
+        code_start = int(layout["code"]["start"], 16)  # type: ignore[index, 
call-overload]
+        assert b"*/" in fixture[code_start:]
+
+        assert read_bundle(TYPESCRIPT_V1_FIXTURE).dag_ids == 
frozenset({"test_dag"})
+
     def test_rejects_metadata_first_legacy_bundle(self, tmp_path):
         payload = _metadata_json("sales")
         (tmp_path / BUNDLE_NAME).write_bytes(METADATA_PREFIX + payload + 
b"\nexport {};\n")
@@ -85,9 +117,9 @@ class TestBundleReader:
         original_header_size = len(bundle.read_bytes().partition(b"\n")[0])
         layout["future_header_field"] = True
         layout["code"]["future_section_field"] = True  # type: ignore[index]
-        # Fixed-width offsets let us relocate both sections after extending 
the header.
+        # Fixed-width offsets let us relocate every section after extending 
the header.
         header_growth = len(LAYOUT_PREFIX) + len(json.dumps(layout).encode()) 
- original_header_size
-        for name in ("metadata", "code"):
+        for name in ("metadata", "source", "code"):
             for field in ("start", "end"):
                 value = int(layout[name][field], 16) + header_growth  # type: 
ignore[index, call-overload]
                 layout[name][field] = f"{value:0{OFFSET_WIDTH}x}"  # type: 
ignore[index]
@@ -168,7 +200,7 @@ class TestBundleReader:
         with pytest.raises(ValueError, match="metadata contains a JavaScript 
line terminator"):
             read_bundle(bundle)
 
-    @pytest.mark.parametrize("section", ["code", "metadata"])
+    @pytest.mark.parametrize("section", ["code", "metadata", "source"])
     def test_requires_every_layout_section(self, tmp_path, section):
         bundle = write_bundle(tmp_path, "sales")
         layout = _read_layout(bundle)
@@ -217,7 +249,9 @@ class TestBundleReader:
         with pytest.raises(ValueError, match="metadata offsets do not match"):
             read_bundle(tmp_path / BUNDLE_NAME)
 
-    @pytest.mark.parametrize("section", [pytest.param("metadata", 
id="metadata-before-decode"), "code"])
+    @pytest.mark.parametrize(
+        "section", [pytest.param("metadata", id="metadata-before-decode"), 
"code", "source"]
+    )
     def test_rejects_section_digest_mismatch(self, tmp_path, section):
         bundle = write_bundle(tmp_path, "sales")
         layout = _read_layout(bundle)
@@ -227,6 +261,85 @@ class TestBundleReader:
         with pytest.raises(ValueError, match=f"{section} SHA-256 mismatch"):
             read_bundle(tmp_path / BUNDLE_NAME)
 
+    def test_rejects_source_offset_mismatch(self, tmp_path):
+        bundle = write_bundle(tmp_path, "sales")
+        layout = _read_layout(bundle)
+        source_start = int(layout["source"]["start"], 16)  # type: 
ignore[index, call-overload]
+        layout["source"]["start"] = f"{source_start + 1:0{OFFSET_WIDTH}x}"  # 
type: ignore[index]
+        _rewrite_layout(bundle, layout)
+
+        with pytest.raises(ValueError, match="source offsets do not match"):
+            read_bundle(tmp_path / BUNDLE_NAME)
+
+    def test_requires_source_comment_immediately_after_metadata(self, 
tmp_path):
+        bundle = write_bundle(tmp_path, "sales")
+        contents = bundle.read_bytes()
+        # Replace the block-comment opener with a line comment of the same 
length,
+        # leaving every declared offset intact.
+        bundle.write_bytes(contents.replace(SOURCE_OPEN, b"//# 
airflowSourc\n", 1))
+
+        with pytest.raises(ValueError, match="no embedded airflow source after 
its metadata"):
+            read_bundle(tmp_path / BUNDLE_NAME)
+
+    def test_rejects_unterminated_source(self, tmp_path):
+        bundle = write_bundle(tmp_path, "sales")
+        contents = bundle.read_bytes()
+        bundle.write_bytes(contents.replace(SOURCE_CLOSE, b"\n#*-\n", 1))
+
+        with pytest.raises(ValueError, match="not closed by its block 
comment"):
+            read_bundle(tmp_path / BUNDLE_NAME)
+
+    def test_rejects_source_declared_past_end_of_file(self, tmp_path):
+        bundle = write_bundle(tmp_path, "sales")
+        layout = _read_layout(bundle)
+        source_end = int(layout["source"]["end"], 16)  # type: ignore[index, 
call-overload]
+        layout["source"]["end"] = f"{source_end + 4096:0{OFFSET_WIDTH}x}"  # 
type: ignore[index]
+        _rewrite_layout(bundle, layout)
+
+        with pytest.raises(ValueError, match="not closed by its block 
comment"):
+            read_bundle(tmp_path / BUNDLE_NAME)
+
+    def test_rejects_oversized_source(self, tmp_path):
+        bundle = write_bundle(tmp_path, "sales")
+        layout = _read_layout(bundle)
+        source_start = int(layout["source"]["start"], 16)  # type: 
ignore[index, call-overload]
+        layout["source"]["end"] = f"{source_start + 1024 * 1024 + 
1:0{OFFSET_WIDTH}x}"  # type: ignore[index]
+        _rewrite_layout(bundle, layout)
+
+        with pytest.raises(ValueError, match="embedded airflow source 
exceeds"):
+            read_bundle(tmp_path / BUNDLE_NAME)
+
+    def test_rejects_unescaped_block_comment_terminator_in_source(self, 
tmp_path):
+        # Node would end the comment at the terminator and execute what 
follows, while both
+        # digests still match. Reject rather than vouch for such a bundle.
+        write_bundle(tmp_path, "sales", source_payload=b'const s = "*/";')
+
+        with pytest.raises(ValueError, match="unescaped block comment 
terminator"):
+            read_bundle(tmp_path / BUNDLE_NAME)
+
+    def test_rejects_source_that_is_not_utf8(self, tmp_path):
+        write_bundle(tmp_path, "sales", source_payload=b"const s = '" + 
bytes([0xFF]) + b"';")
+
+        with pytest.raises(ValueError, match="embedded airflow source is not 
valid UTF-8"):
+            read_bundle_source(tmp_path / BUNDLE_NAME)
+
+    @pytest.mark.parametrize(
+        "source",
+        [
+            pytest.param(b"export {};" + b"\n", id="plain"),
+            pytest.param(b"/** doc */" + b"\nexport {};\n", id="doc-comment"),
+            pytest.param(b'const s = "*/";\n', id="terminator"),
+            pytest.param(b'const s = "*\\/";' + b"\n", id="escaped-slash"),
+            pytest.param(b'const s = "*\\\\";' + b"\n", id="double-backslash"),
+            pytest.param('const s = "\u65e5\u672c\u8a9e";\n'.encode(), 
id="non-ascii"),
+            pytest.param(b"line one\nline two\nline three\n", id="multi-line"),
+        ],
+    )
+    def test_round_trips_embedded_source_exactly(self, tmp_path, source):
+        bundle = write_bundle(tmp_path, "sales", source=source)
+
+        assert read_bundle_source(bundle).encode() == source
+
     def test_rejects_truncated_code(self, tmp_path):
         bundle = write_bundle(tmp_path, "sales")
         bundle.write_bytes(bundle.read_bytes()[:-1])
@@ -264,16 +377,16 @@ class TestBundleReader:
         with pytest.raises(ValueError, match="embedded airflow metadata must 
contain a mapping"):
             read_bundle(tmp_path / BUNDLE_NAME)
 
-    @pytest.mark.parametrize("dags", [None, []], ids=["missing", 
"not-a-mapping"])
-    def test_rejects_missing_or_malformed_dags(self, tmp_path, dags):
+    @pytest.mark.parametrize("task_handlers", [None, []], ids=["missing", 
"not-a-mapping"])
+    def test_rejects_missing_or_malformed_task_handlers(self, tmp_path, 
task_handlers):
         metadata = json.loads(_metadata_json("sales"))
-        if dags is None:
-            del metadata["dags"]
+        if task_handlers is None:
+            del metadata["task_handlers"]
         else:
-            metadata["dags"] = dags
+            metadata["task_handlers"] = task_handlers
         write_bundle(tmp_path, metadata_payload=json.dumps(metadata).encode())
 
-        with pytest.raises(ValueError, match="metadata must contain a dags 
mapping"):
+        with pytest.raises(ValueError, match="metadata must contain a 
task_handlers mapping"):
             read_bundle(tmp_path / BUNDLE_NAME)
 
     def test_rejects_oversized_metadata(self, tmp_path):
@@ -363,7 +476,8 @@ class TestBundleReader:
         read_bundle(tmp_path / BUNDLE_NAME)
         read_bundle(tmp_path / BUNDLE_NAME)
 
-        assert hash_region.call_count == 2
+        # Three regions hashed on the first read, none on the second.
+        assert hash_region.call_count == 3
 
     @mock.patch.object(_reader.os, "fstat", autospec=True)
     def test_cache_uses_ctime_to_detect_corruption_with_restored_mtime(self, 
fstat, tmp_path):
@@ -389,12 +503,13 @@ class TestBundleReader:
     def test_digest_cache_evicts_least_recently_used_entry(self):
         cache = _reader._BundleDigestCache(maxsize=2)
         section = _reader._DeclaredSection(start=0, end=1, sha256=b"0" * 32)
-        digests = _reader._ComputedDigests(metadata=b"1" * 32, code=b"2" * 32)
+        digests = _reader._ComputedDigests(metadata=b"1" * 32, source=b"2" * 
32, code=b"3" * 32)
 
         def build_key(inode):
             return _reader._DigestCacheKey(
                 path="bundle.min.mjs",
                 metadata=section,
+                source=section,
                 code=section,
                 device=1,
                 inode=inode,
diff --git a/ts-sdk/.prettierignore b/ts-sdk/.prettierignore
index 4d0d3a67c0e..340b10cba65 100644
--- a/ts-sdk/.prettierignore
+++ b/ts-sdk/.prettierignore
@@ -18,3 +18,6 @@
 dist/
 coverage/
 node_modules/
+
+# Byte-exact bundle artifacts: reformatting one invalidates its embedded 
digests.
+tests/cli/fixtures/*.min.mjs
diff --git a/ts-sdk/README.md b/ts-sdk/README.md
index e82dc090a05..e2aba6595f9 100644
--- a/ts-sdk/README.md
+++ b/ts-sdk/README.md
@@ -266,9 +266,11 @@ airflow-ts-pack src/main.ts --outdir dist
 It bundles the entrypoint into a minified `dist/bundle.min.mjs` with esbuild, 
then runs that bundle with
 `--airflow-metadata` so it reports its own registered Dag/task pairs and 
supervisor schema version. The manifest is
 embedded as a compact JSON `//# airflowMetadata=...` comment after a leading 
compact JSON `//# airflowBundle=...`
-layout descriptor. The CLI records the integrity metadata for both regions in 
that descriptor, so a coordinator that
-is handed a bundle whose content was replaced fails loudly instead of running 
it. The result is one deployable file
-with no hand-written metadata sidecar.
+layout descriptor, and the entry module is embedded verbatim in a `/*# 
airflowSource ... #*/` block comment so
+Airflow can show the source a bundle was authored from, which its minified 
code no longer is. The CLI records the
+integrity metadata for all three regions in that descriptor, so a coordinator 
that is handed a bundle whose content
+was replaced fails loudly instead of running it. The result is one deployable 
file with no hand-written metadata
+sidecar.
 
 Pass `--outfile <path>` instead of `--outdir` to name the artifact yourself, 
so one bundle directory can hold several
 bundles. The name must still end in `.min.mjs`, which is how `NodeCoordinator` 
finds bundles.
diff --git a/ts-sdk/src/cli/bundle-encoder.ts b/ts-sdk/src/cli/bundle-encoder.ts
index 5870b0faf84..f712ae33521 100644
--- a/ts-sdk/src/cli/bundle-encoder.ts
+++ b/ts-sdk/src/cli/bundle-encoder.ts
@@ -25,14 +25,15 @@
  *
  *   airflowBundle header
  *   -> airflowMetadata
+ *   -> airflowSource
  *   -> executable JavaScript
  *
- * The header tells Airflow where each region begins and ends and carries the
- * digest used to verify each one. Metadata describes what the bundle can 
serve,
- * and executable JavaScript runs its task handlers.
+ * The header records each region's byte range and digest. Metadata describes 
what the bundle can
+ * serve, the source region carries the entrypoint as written, and the 
executable JavaScript runs
+ * the task handlers.
  *
- * This module owns the on-disk encoding. Readers must use the header's named
- * byte ranges rather than relying on incidental line positions.
+ * This module owns the on-disk encoding. Readers must use the header's named 
byte ranges rather
+ * than incidental line positions.
  */
 
 import { createHash } from "node:crypto";
@@ -41,15 +42,20 @@ import type { BundleManifest } from 
"../coordinator/manifest.js";
 
 const AIRFLOW_BUNDLE_METADATA_VERSION = "1.0";
 const EMBEDDED_METADATA_MAX_BYTES = 1024 * 1024;
+const EMBEDDED_SOURCE_MAX_BYTES = 1024 * 1024;
 const OFFSET_HEX_WIDTH = 16;
 
 export const EMBEDDED_METADATA_PREFIX = "//# airflowMetadata=";
 export const EMBEDDED_LAYOUT_PREFIX = "//# airflowBundle=";
+/** A block comment, because the entrypoint spans more than one line. */
+export const EMBEDDED_SOURCE_OPEN = "/*# airflowSource\n";
+export const EMBEDDED_SOURCE_CLOSE = "\n#*/\n";
 
 export interface BundleEncoderInput {
   bundleManifest: BundleManifest;
   sdkVersion: string;
   entrypointName: string;
+  entrypointSource: string;
   executable: Uint8Array;
 }
 
@@ -57,7 +63,7 @@ interface BundleMetadata {
   airflow_bundle_metadata_version: string;
   sdk: { language: string; version: string; supervisor_schema_version: string 
};
   source: string;
-  dags: BundleManifest["dags"];
+  task_handlers: BundleManifest["task_handlers"];
 }
 
 interface VerifiedByteRange {
@@ -69,33 +75,45 @@ interface VerifiedByteRange {
 interface BundleHeader {
   code: VerifiedByteRange;
   metadata: VerifiedByteRange;
+  source: VerifiedByteRange;
 }
 
 export function encodeBundle(input: BundleEncoderInput): Buffer {
   const metadata = encodeMetadata(input);
+  const source = encodeSource(input.entrypointSource);
   const executable = encodeExecutable(input.executable);
-  const header = encodeHeader({ metadata, executable });
+  const header = encodeHeader({ metadata, source, executable });
 
-  return Buffer.concat([header, metadata, executable]);
+  return Buffer.concat([header, metadata, source, executable]);
 }
 
-function encodeHeader(regions: { metadata: Buffer; executable: Buffer }): 
Buffer {
+function encodeHeader(regions: { metadata: Buffer; source: Buffer; executable: 
Buffer }): Buffer {
+  // Each digest covers the payload only. The framing markers and newlines are 
re-derived.
   const metadataPayload = regions.metadata.subarray(
     Buffer.byteLength(EMBEDDED_METADATA_PREFIX),
     -1,
   );
+  const sourcePayload = regions.source.subarray(
+    Buffer.byteLength(EMBEDDED_SOURCE_OPEN),
+    -Buffer.byteLength(EMBEDDED_SOURCE_CLOSE),
+  );
   const digests = {
     code: computeSha256(regions.executable),
     metadata: computeSha256(metadataPayload),
+    source: computeSha256(sourcePayload),
   };
   const zeroOffset = "0".repeat(OFFSET_HEX_WIDTH);
   const placeholderHeader = renderHeader({
     code: { start: zeroOffset, end: zeroOffset, sha256: digests.code },
     metadata: { start: zeroOffset, end: zeroOffset, sha256: digests.metadata },
+    source: { start: zeroOffset, end: zeroOffset, sha256: digests.source },
   });
   const metadataStart = placeholderHeader.length + 
Buffer.byteLength(EMBEDDED_METADATA_PREFIX);
   const metadataEnd = metadataStart + metadataPayload.length;
-  const codeStart = placeholderHeader.length + regions.metadata.length;
+  const sourceStart =
+    placeholderHeader.length + regions.metadata.length + 
Buffer.byteLength(EMBEDDED_SOURCE_OPEN);
+  const sourceEnd = sourceStart + sourcePayload.length;
+  const codeStart = placeholderHeader.length + regions.metadata.length + 
regions.source.length;
   const codeEnd = codeStart + regions.executable.length;
   const header = renderHeader({
     code: {
@@ -108,6 +126,11 @@ function encodeHeader(regions: { metadata: Buffer; 
executable: Buffer }): Buffer
       end: formatOffset(metadataEnd),
       sha256: digests.metadata,
     },
+    source: {
+      start: formatOffset(sourceStart),
+      end: formatOffset(sourceEnd),
+      sha256: digests.source,
+    },
   });
   if (header.length !== placeholderHeader.length) {
     throw new Error("Bundle header changed length while resolving section 
offsets");
@@ -115,6 +138,33 @@ function encodeHeader(regions: { metadata: Buffer; 
executable: Buffer }): Buffer
   return header;
 }
 
+/**
+ * Wrap the entrypoint as written in a block comment.
+ *
+ * A comment terminator would splice the rest of the payload into executable 
position, so it is
+ * escaped. `*\\` is escaped too, which keeps the transformation reversible.
+ */
+function encodeSource(entrypointSource: string): Buffer {
+  const payload = Buffer.from(escapeBlockComment(entrypointSource), "utf-8");
+  const source = Buffer.concat([
+    Buffer.from(EMBEDDED_SOURCE_OPEN, "ascii"),
+    payload,
+    Buffer.from(EMBEDDED_SOURCE_CLOSE, "ascii"),
+  ]);
+  if (source.length > EMBEDDED_SOURCE_MAX_BYTES) {
+    throw new Error(
+      `Embedded entrypoint source is ${source.length} bytes, ` +
+        `over the ${EMBEDDED_SOURCE_MAX_BYTES} byte limit; move code out of 
the entrypoint ` +
+        `into imported modules`,
+    );
+  }
+  return source;
+}
+
+function escapeBlockComment(source: string): string {
+  return source.replaceAll(/\*([\\/])/g, "*\\$1");
+}
+
 function encodeMetadata(input: BundleEncoderInput): Buffer {
   const payload = Buffer.from(
     JSON.stringify(buildBundleMetadata(input))
@@ -152,7 +202,7 @@ function buildBundleMetadata(input: BundleEncoderInput): 
BundleMetadata {
       supervisor_schema_version: 
input.bundleManifest.supervisor_schema_version,
     },
     source: input.entrypointName,
-    dags: input.bundleManifest.dags,
+    task_handlers: input.bundleManifest.task_handlers,
   };
 }
 
diff --git a/ts-sdk/src/cli/pack.ts b/ts-sdk/src/cli/pack.ts
index 0403d8047bd..10e8c5fbbdc 100644
--- a/ts-sdk/src/cli/pack.ts
+++ b/ts-sdk/src/cli/pack.ts
@@ -18,7 +18,8 @@
  */
 
 // airflow-ts-pack: bundle a TypeScript entrypoint into the single artifact 
NodeCoordinator consumes.
-// `bundle.min.mjs` carries the metadata and an integrity layout descriptor in 
JavaScript comments.
+// `bundle.min.mjs` carries the metadata, the entrypoint source, and an 
integrity layout descriptor
+// in JavaScript comments.
 //
 // Build first, then run the built bundle with --airflow-metadata so the 
manifest comes from the
 // bundle's own Dag registry and schema version, never from a hand-written 
sidecar.
@@ -46,7 +47,8 @@ const MANIFEST_MAX_BUFFER_BYTES = 64 * 1024 * 1024;
 const USAGE = `Usage: airflow-ts-pack <entry> [--outdir <dir> | --outfile 
<path>] [--source <name>]
 
 Bundles <entry> into a minified ${BUNDLE_FILENAME} with esbuild and embeds the
-airflow metadata generated from the bundle's served Dags.
+airflow metadata generated from the bundle's served Dags, plus <entry> itself 
as
+the readable source Airflow displays for the bundle.
 
 Options:
   --outdir <dir>    Output directory, holding ${BUNDLE_FILENAME} (default: 
dist)
@@ -156,7 +158,7 @@ function readBundleManifest(bundlePath: string): 
BundleManifest {
   const manifest = parsed;
   // The line is whatever the bundle printed and nothing downstream 
re-validates
   // it, so check each Dag entry down to the task-id element.
-  for (const [dagId, dag] of Object.entries(manifest.dags)) {
+  for (const [dagId, dag] of Object.entries(manifest.task_handlers)) {
     if (dag == null || !isTaskIdList(dag.tasks)) {
       throw new Error(
         `Bundle produced ${AIRFLOW_METADATA_FLAG} output with a malformed 
entry for Dag "${dagId}"`,
@@ -171,16 +173,17 @@ function readBundleManifest(bundlePath: string): 
BundleManifest {
 // raw TypeError rather than a report about the bundle.
 function isBundleManifest(value: unknown): value is BundleManifest {
   if (typeof value !== "object" || value === null || Array.isArray(value)) 
return false;
-  const { supervisor_schema_version: version, dags } = value as 
Partial<BundleManifest>;
+  const { supervisor_schema_version: version, task_handlers: taskHandlers } =
+    value as Partial<BundleManifest>;
   return (
     // Rendered into the manifest verbatim, where the schema requires a 
non-empty
     // string, so a truthy number or boolean would travel to Airflow as-is.
     typeof version === "string" &&
     version.length > 0 &&
-    typeof dags === "object" &&
-    dags !== null &&
+    typeof taskHandlers === "object" &&
+    taskHandlers !== null &&
     // An array would pass the typeof check and yield Dags named "0", "1", ...
-    !Array.isArray(dags)
+    !Array.isArray(taskHandlers)
   );
 }
 
@@ -219,7 +222,7 @@ export async function runPack(argv: readonly string[]): 
Promise<void> {
     });
 
     const manifest = readBundleManifest(stagingPath);
-    const dagEntries = Object.entries(manifest.dags);
+    const dagEntries = Object.entries(manifest.task_handlers);
     if (dagEntries.length === 0) {
       throw new Error(
         `${args.entry} served nothing; register Dags or task handlers with 
bundle.register(...)`,
@@ -232,12 +235,14 @@ export async function runPack(argv: readonly string[]): 
Promise<void> {
         process.stderr.write(`warning: dag ${JSON.stringify(dagId)} has no 
tasks\n`);
       }
     }
-    warnOnSuspiciousIds(manifest.dags);
+    warnOnSuspiciousIds(manifest.task_handlers);
 
     const bundle = encodeBundle({
       bundleManifest: manifest,
       sdkVersion: readSdkVersion(),
       entrypointName: args.source,
+      // Only the entrypoint, like the Java SDK's Airflow-Java-SDK-Dag-Code 
attribute.
+      entrypointSource: readFileSync(args.entry, "utf-8"),
       executable: readFileSync(stagingPath),
     });
     writeFileSync(bundlePath, bundle);
@@ -245,5 +250,5 @@ export async function runPack(argv: readonly string[]): 
Promise<void> {
     rmSync(stagingPath, { force: true });
   }
 
-  console.log(`Wrote ${bundlePath} (airflow metadata and integrity embedded)`);
+  console.log(`Wrote ${bundlePath} (airflow metadata, source, and integrity 
embedded)`);
 }
diff --git a/ts-sdk/src/cli/validate.ts b/ts-sdk/src/cli/validate.ts
index 91fc6e35482..1d8098b3fcf 100644
--- a/ts-sdk/src/cli/validate.ts
+++ b/ts-sdk/src/cli/validate.ts
@@ -31,12 +31,12 @@ type WarnFn = (message: string) => void;
  * one depend on server configuration the packer cannot see.
  */
 export function warnOnSuspiciousIds(
-  dags: BundleManifest["dags"],
+  taskHandlers: BundleManifest["task_handlers"],
   warn: WarnFn = (message) => process.stderr.write(`${message}\n`),
 ): void {
-  for (const dagId of Object.keys(dags).sort()) {
+  for (const dagId of Object.keys(taskHandlers).sort()) {
     warnOnSuspiciousId(`dag id ${JSON.stringify(dagId)}`, dagId, warn);
-    for (const taskId of dags[dagId]!.tasks) {
+    for (const taskId of taskHandlers[dagId]!.tasks) {
       warnOnSuspiciousId(
         `task id ${JSON.stringify(taskId)} in dag ${JSON.stringify(dagId)}`,
         taskId,
diff --git a/ts-sdk/src/coordinator/manifest.ts 
b/ts-sdk/src/coordinator/manifest.ts
index 3a35fc46e09..97661257efd 100644
--- a/ts-sdk/src/coordinator/manifest.ts
+++ b/ts-sdk/src/coordinator/manifest.ts
@@ -25,23 +25,22 @@ export const AIRFLOW_METADATA_FLAG = "--airflow-metadata";
 /** Marks the manifest line on stdout, which import-time logging may also 
reach. */
 export const AIRFLOW_METADATA_SENTINEL = "__AIRFLOW_METADATA__ ";
 
-/** Bundle manifest fields only the built bundle itself knows: the schema
- *  version it was compiled against and the Dag/task pairs it registered.
- *  Registered Dags without tasks appear with an empty `tasks` list so
- *  `airflow-ts-pack` (which runs `node bundle.mjs --airflow-metadata` to
- *  read this) can warn about them instead of silently dropping them. */
+/** Bundle manifest fields only the built bundle itself knows: the schema 
version it was compiled
+ *  against, and the task handlers it registered grouped by Dag. Named 
`task_handlers` because a
+ *  TypeScript bundle provides handlers for Dags declared elsewhere, not Dag 
definitions. A Dag with
+ *  no handlers keeps an empty `tasks` list so `airflow-ts-pack` can warn 
instead of dropping it. */
 export interface BundleManifest {
   supervisor_schema_version: string;
-  dags: Record<string, { tasks: string[] }>;
+  task_handlers: Record<string, { tasks: string[] }>;
 }
 
 export function buildBundleManifest(bundle: Bundle): BundleManifest {
-  const dags: BundleManifest["dags"] = {};
+  const taskHandlers: BundleManifest["task_handlers"] = {};
   for (const { dagId, tasks } of listBundleDags(bundle)) {
     if (typeof dagId !== "string") {
       throw new Error("Dag ID must be a string");
     }
-    Object.defineProperty(dags, dagId, {
+    Object.defineProperty(taskHandlers, dagId, {
       configurable: true,
       enumerable: true,
       value: { tasks: [...tasks] },
@@ -50,6 +49,6 @@ export function buildBundleManifest(bundle: Bundle): 
BundleManifest {
   }
   return {
     supervisor_schema_version: SUPERVISOR_API_VERSION,
-    dags,
+    task_handlers: taskHandlers,
   };
 }
diff --git a/ts-sdk/tests/cli/fixtures/bundle-v1.min.mjs 
b/ts-sdk/tests/cli/fixtures/bundle-v1.min.mjs
new file mode 100644
index 00000000000..91115f3ba89
--- /dev/null
+++ b/ts-sdk/tests/cli/fixtures/bundle-v1.min.mjs
@@ -0,0 +1,18 @@
+//# 
airflowBundle={"code":{"start":"000000000000039a","end":"00000000000004a9","sha256":"fa1afadd7cb7147d54e4870b849dd2390e5ad75f6153d040e7381ae8a76a5a29"},"metadata":{"start":"00000000000001c9","end":"0000000000000296","sha256":"93ef7891d3e562383d4fd95332d9cf3b0662e581623dff811443350d16a84479"},"source":{"start":"00000000000002a9","end":"0000000000000395","sha256":"3e7a601893ae2cc92dd72dcc596d3fbfa19cafbaa0d3f5768973c49c3daee1c9"}}
+//# 
airflowMetadata={"airflow_bundle_metadata_version":"1.0","sdk":{"language":"typescript","version":"0.1.0","supervisor_schema_version":"2026-06-16"},"source":"entry.ts","task_handlers":{"test_dag":{"tasks":["test_task"]}}}
+/*# airflowSource
+/** Handlers for the test Dag. *\/
+import { Bundle, Dag } from "apache-airflow-ts-sdk";
+
+const TERMINATOR = /\*\\//;
+const dag = new Dag("test_dag");
+dag.task("test_task", async () => TERMINATOR.source);
+
+await new Bundle(dag).serve();
+
+#*/
+var e=async function(){return"extracted"};await e();
+/*! 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.
+ */
diff --git a/ts-sdk/tests/cli/fixtures/bundle-v1.mjs 
b/ts-sdk/tests/cli/fixtures/bundle-v1.mjs
deleted file mode 100644
index 9edb3b8d3c2..00000000000
--- a/ts-sdk/tests/cli/fixtures/bundle-v1.mjs
+++ /dev/null
@@ -1,24 +0,0 @@
-//# 
airflowBundle={"code":{"start":"0000000000000203","end":"000000000000057a","sha256":"bd26f9a6295069aef9eff45213377a49a448dcc0973da182a29ffb927f259e8e"},"metadata":{"start":"000000000000013e","end":"0000000000000202","sha256":"a51dfd6f0c9e8ea867900e55c0387b556d3cb0e98321b62d4625f522ed465041"}}
-//# 
airflowMetadata={"airflow_bundle_metadata_version":"1.0","sdk":{"language":"typescript","version":"0.1.0","supervisor_schema_version":"2026-06-16"},"source":"entry.ts","dags":{"test_dag":{"tasks":["test_task"]}}}
-/*!
- * 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.
- */
-
-import { Bundle } from "../../../src/index.js";
-
-await new Bundle().serve();
diff --git a/ts-sdk/tests/cli/pack.test.ts b/ts-sdk/tests/cli/pack.test.ts
index fb45586a254..fad4d50a8a8 100644
--- a/ts-sdk/tests/cli/pack.test.ts
+++ b/ts-sdk/tests/cli/pack.test.ts
@@ -28,6 +28,8 @@ import { afterEach, describe, expect, it, vi } from "vitest";
 import {
   EMBEDDED_LAYOUT_PREFIX,
   EMBEDDED_METADATA_PREFIX,
+  EMBEDDED_SOURCE_CLOSE,
+  EMBEDDED_SOURCE_OPEN,
   encodeBundle,
 } from "../../src/cli/bundle-encoder.js";
 import { parsePackArgs, runPack } from "../../src/cli/pack.js";
@@ -35,10 +37,40 @@ import { SUPERVISOR_API_VERSION } from 
"../../src/coordinator/protocol.js";
 import { AIRFLOW_METADATA_SENTINEL } from "../../src/coordinator/manifest.js";
 
 const FIXTURE_ENTRY = fileURLToPath(new URL("fixtures/entry.ts", 
import.meta.url));
-const GOLDEN_BUNDLE = fileURLToPath(new URL("fixtures/bundle-v1.mjs", 
import.meta.url));
+const GOLDEN_BUNDLE = fileURLToPath(new URL("fixtures/bundle-v1.min.mjs", 
import.meta.url));
 const NOISY_ENTRY = fileURLToPath(new URL("fixtures/noisy-entry.ts", 
import.meta.url));
 const EMPTY_ENTRY = fileURLToPath(new URL("fixtures/empty-entry.ts", 
import.meta.url));
 const SDK_INDEX = fileURLToPath(new URL("../../src/index.ts", 
import.meta.url));
+
+// Shaped like what airflow-ts-pack ships: a code region nobody could read, 
and a source region
+// that needs escaping. Real esbuild output is not used here because the 
assertion is byte-exact
+// and would churn on every esbuild bump. The runPack tests cover real output.
+const GOLDEN_CODE = Buffer.from(
+  [
+    'var e=async function(){return"extracted"};await e();',
+    // esbuild relocates dependencies' banners to the end, putting a comment 
terminator in the
+    // code region. Only the source region may not hold one, and the fixture 
pins that.
+    "/*! 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.",
+    " */",
+    "",
+  ].join("\n"),
+);
+
+// Carries both escape branches: the doc comment ends in a terminator, and the 
regex holds a star
+// followed by a backslash. The Python reader is then validated against real 
encoder output.
+const GOLDEN_SOURCE = [
+  "/** Handlers for the test Dag. */",
+  'import { Bundle, Dag } from "apache-airflow-ts-sdk";',
+  "",
+  "const TERMINATOR = /\\*\\//;",
+  'const dag = new Dag("test_dag");',
+  'dag.task("test_task", async () => TERMINATOR.source);',
+  "",
+  "await new Bundle(dag).serve();",
+  "",
+].join("\n");
 const SDK_VERSION = (
   JSON.parse(readFileSync(new URL("../../package.json", import.meta.url), 
"utf-8")) as {
     version: string;
@@ -48,6 +80,7 @@ const SDK_VERSION = (
 interface TestBundleHeader {
   code: { end: string; sha256: string; start: string };
   metadata: { end: string; sha256: string; start: string };
+  source: { end: string; sha256: string; start: string };
 }
 
 function parseHeader(line: string): TestBundleHeader {
@@ -104,10 +137,11 @@ describe("encodeBundle", () => {
     const bundle = encodeBundle({
       bundleManifest: {
         supervisor_schema_version: "2026-06-16",
-        dags: { my_dag: { tasks: ["a", 'b"c'] } },
+        task_handlers: { my_dag: { tasks: ["a", 'b"c'] } },
       },
       sdkVersion: "0.1.0",
       entrypointName: 'we"ird.ts',
+      entrypointSource: "const x = 1;\nconsole.log(x);\n",
       executable,
     });
 
@@ -116,40 +150,90 @@ describe("encodeBundle", () => {
     const offset = (value: string): number => Number.parseInt(value, 16);
     const metadataStart = offset(header.metadata.start);
     const metadataEnd = offset(header.metadata.end);
+    const sourceStart = offset(header.source.start);
+    const sourceEnd = offset(header.source.end);
     const executableStart = offset(header.code.start);
     const executableEnd = offset(header.code.end);
 
     expect(metadataStart).toBe(firstNewline + 1 + 
Buffer.byteLength(EMBEDDED_METADATA_PREFIX));
-    expect(executableStart).toBe(metadataEnd + 1);
+    expect(sourceStart).toBe(metadataEnd + 1 + 
Buffer.byteLength(EMBEDDED_SOURCE_OPEN));
+    expect(executableStart).toBe(sourceEnd + 
Buffer.byteLength(EMBEDDED_SOURCE_CLOSE));
     expect(executableEnd).toBe(bundle.length);
     expect(bundle.subarray(executableStart, 
executableEnd)).toEqual(executable);
+    expect(bundle.subarray(sourceStart, sourceEnd).toString("utf-8")).toBe(
+      "const x = 1;\nconsole.log(x);\n",
+    );
 
     const metadata = bundle.subarray(metadataStart, 
metadataEnd).toString("utf-8");
     expect(metadata).toBe(
-      
'{"airflow_bundle_metadata_version":"1.0","sdk":{"language":"typescript","version":"0.1.0","supervisor_schema_version":"2026-06-16"},"source":"we\\"ird.ts","dags":{"my_dag":{"tasks":["a","b\\"c"]}}}',
+      
'{"airflow_bundle_metadata_version":"1.0","sdk":{"language":"typescript","version":"0.1.0","supervisor_schema_version":"2026-06-16"},"source":"we\\"ird.ts","task_handlers":{"my_dag":{"tasks":["a","b\\"c"]}}}',
     );
 
-    expect(header).not.toHaveProperty("source");
     expect(header).not.toHaveProperty("version");
-    expect(bundle.toString("utf-8")).not.toContain("airflowSource");
+  });
+
+  it.each([
+    { label: "a block comment terminator", source: "/** doc */\nexport {};\n" 
},
+    { label: "an escaped slash after a star", source: 'const s = "*\\\\/";\n' 
},
+    { label: "a double backslash after a star", source: 'const s = 
"*\\\\\\\\";\n' },
+  ])("escapes $label so the source region cannot close early", ({ source }) => 
{
+    const bundle = encodeBundle({
+      bundleManifest: { supervisor_schema_version: "2026-06-16", 
task_handlers: {} },
+      sdkVersion: "0.1.0",
+      entrypointName: "entry.ts",
+      entrypointSource: source,
+      executable: Buffer.from("export {};\n"),
+    });
+    const header = parseHeader(bundle.subarray(0, 
bundle.indexOf("\n")).toString("ascii"));
+    const payload = bundle
+      .subarray(Number.parseInt(header.source.start, 16), 
Number.parseInt(header.source.end, 16))
+      .toString("utf-8");
+
+    // Anything else would end the comment where Node reads the file.
+    expect(payload).not.toContain("*/");
+    // Reversing the escaping recovers the entrypoint byte for byte.
+    expect(payload.replaceAll(/\*\\([\\/])/g, "*$1")).toBe(source);
+  });
+
+  it("rejects an entrypoint source over the embedded size limit", () => {
+    expect(() =>
+      encodeBundle({
+        bundleManifest: { supervisor_schema_version: "2026-06-16", 
task_handlers: {} },
+        sdkVersion: "0.1.0",
+        entrypointName: "entry.ts",
+        entrypointSource: "x".repeat(1024 * 1024 + 1),
+        executable: Buffer.from("export {};\n"),
+      }),
+    ).toThrow("over the 1048576 byte limit");
   });
 
   it("matches the golden bundle", () => {
-    const executable = readFileSync(EMPTY_ENTRY);
     const bundle = encodeBundle({
       bundleManifest: {
         supervisor_schema_version: "2026-06-16",
-        dags: { test_dag: { tasks: ["test_task"] } },
+        task_handlers: { test_dag: { tasks: ["test_task"] } },
       },
       sdkVersion: "0.1.0",
       entrypointName: "entry.ts",
-      executable,
+      entrypointSource: GOLDEN_SOURCE,
+      executable: GOLDEN_CODE,
     });
 
     expect(bundle).toEqual(readFileSync(GOLDEN_BUNDLE));
     const firstNewline = bundle.indexOf("\n");
     const header = parseHeader(bundle.subarray(0, 
firstNewline).toString("ascii"));
-    for (const section of [header.metadata, header.code]) {
+    const offset = (value: string): number => Number.parseInt(value, 16);
+    const source = bundle.subarray(offset(header.source.start), 
offset(header.source.end));
+    // Stored escaped and reversible: the byte-level agreement the Python 
reader is checked against.
+    expect(source.toString("utf-8")).not.toContain("*/");
+    expect(source.toString("utf-8")).not.toBe(GOLDEN_SOURCE);
+    expect(source.toString("utf-8").replaceAll(/\*\\([\\/])/g, 
"*$1")).toBe(GOLDEN_SOURCE);
+    // A terminator in the code region is fine. Only the source comment cannot 
hold one.
+    expect(
+      bundle.subarray(offset(header.code.start), 
offset(header.code.end)).toString("utf-8"),
+    ).toContain("*/");
+
+    for (const section of [header.metadata, header.source, header.code]) {
       expect(section.start).toMatch(/^[0-9a-f]{16}$/);
       expect(section.end).toMatch(/^[0-9a-f]{16}$/);
     }
@@ -159,10 +243,11 @@ describe("encodeBundle", () => {
     const bundle = encodeBundle({
       bundleManifest: {
         supervisor_schema_version: "2026-06-16",
-        dags: { "line\u2028separator": { tasks: ["paragraph\u2029separator"] } 
},
+        task_handlers: { "line\u2028separator": { tasks: 
["paragraph\u2029separator"] } },
       },
       sdkVersion: "0.1.0",
       entrypointName: "entry.ts",
+      entrypointSource: "export {};\n",
       executable: Buffer.from("export {};\n"),
     });
     const metadataLine = bundle.toString("utf-8").split("\n")[1]!;
@@ -172,7 +257,7 @@ describe("encodeBundle", () => {
     expect(metadataLine).toContain("\\u2028");
     expect(metadataLine).toContain("\\u2029");
     
expect(JSON.parse(metadataLine.slice(EMBEDDED_METADATA_PREFIX.length))).toHaveProperty(
-      "dags.line\u2028separator.tasks",
+      "task_handlers.line\u2028separator.tasks",
       ["paragraph\u2029separator"],
     );
   });
@@ -225,7 +310,7 @@ describe("runPack", () => {
         supervisor_schema_version: SUPERVISOR_API_VERSION,
       },
       source: "entry.ts",
-      dags: {
+      task_handlers: {
         fixture_dag: { tasks: ["extract", "transform"] },
         other_dag: { tasks: ["solo"] },
       },
@@ -261,9 +346,16 @@ describe("runPack", () => {
       offset(layout.metadata.end),
     );
     
expect(createHash("sha256").update(metadataPayload).digest("hex")).toBe(layout.metadata.sha256);
-    
expect(JSON.parse(metadataPayload.toString("utf-8"))).toHaveProperty("dags.fixture_dag");
-    expect(layout).not.toHaveProperty("source");
-    expect(bundle.toString("utf-8")).not.toContain("airflowSource");
+    expect(JSON.parse(metadataPayload.toString("utf-8"))).toHaveProperty(
+      "task_handlers.fixture_dag",
+    );
+
+    const source = bundle.subarray(offset(layout.source.start), 
offset(layout.source.end));
+    
expect(createHash("sha256").update(source).digest("hex")).toBe(layout.source.sha256);
+    // Stored escaped, because the fixture's own license header ends in a 
comment terminator.
+    expect(source.toString("utf-8").replaceAll(/\*\\([\\/])/g, "*$1")).toBe(
+      readFileSync(FIXTURE_ENTRY, "utf-8"),
+    );
   });
 
   it("minifies the code region and keeps it runnable", async () => {
@@ -303,7 +395,7 @@ describe("runPack", () => {
     expect(readFileSync(target).subarray(0, 
EMBEDDED_LAYOUT_PREFIX.length).toString()).toBe(
       EMBEDDED_LAYOUT_PREFIX,
     );
-    
expect(JSON.parse(readEmbeddedMetadata(target))).toHaveProperty("dags.fixture_dag");
+    
expect(JSON.parse(readEmbeddedMetadata(target))).toHaveProperty("task_handlers.fixture_dag");
   });
 
   it("keeps a shebang entry runnable and reads the manifest past import-time 
logging", async () => {
@@ -313,12 +405,16 @@ describe("runPack", () => {
     const bundlePath = path.join(outdir, "bundle.min.mjs");
     const bundle = readFileSync(bundlePath, "utf-8");
     expect(bundle.startsWith(EMBEDDED_LAYOUT_PREFIX)).toBe(true);
-    expect(bundle).not.toContain("#!/usr/bin/env node");
+    // The entrypoint is embedded as written and opens with a shebang, so 
check only the code region.
+    const codeRegion = readFileSync(bundlePath).subarray(
+      Number.parseInt(parseHeader(bundle.split("\n")[0]!).code.start, 16),
+    );
+    expect(codeRegion.toString("utf-8")).not.toContain("#!/usr/bin/env node");
     expect(existsSync(path.join(outdir, 
"bundle.pack-staging.mjs"))).toBe(false);
 
     const metadataLine = bundle.split("\n")[1]!;
     const metadata = 
JSON.parse(metadataLine.slice(EMBEDDED_METADATA_PREFIX.length));
-    expect(metadata).toHaveProperty("dags.noisy_dag");
+    expect(metadata).toHaveProperty("task_handlers.noisy_dag");
 
     execFileSync(process.execPath, [bundlePath, "--airflow-metadata"], { 
encoding: "utf-8" });
   });
@@ -406,22 +502,25 @@ describe("runPack", () => {
 
   // A bundle can print the sentinel itself, so nothing on that line is 
trusted.
   it.each([
-    ['{ supervisor_schema_version: "1", dags: { broken_dag: {} } }', 
"malformed entry"],
+    ['{ supervisor_schema_version: "1", task_handlers: { broken_dag: {} } }', 
"malformed entry"],
     [
-      '{ supervisor_schema_version: "1", dags: { broken_dag: { tasks: ["ok", 
7] } } }',
+      '{ supervisor_schema_version: "1", task_handlers: { broken_dag: { tasks: 
["ok", 7] } } }',
       "malformed entry",
     ],
     [
-      '{ supervisor_schema_version: "1", dags: { broken_dag: { tasks: [""] } } 
}',
+      '{ supervisor_schema_version: "1", task_handlers: { broken_dag: { tasks: 
[""] } } }',
       "malformed entry",
     ],
-    ['{ supervisor_schema_version: "1", dags: [{ tasks: ["a"] }] }', 
"incomplete"],
+    ['{ supervisor_schema_version: "1", task_handlers: [{ tasks: ["a"] }] }', 
"incomplete"],
     // Was read off before the document itself was checked, so it surfaced as a
     // raw TypeError.
     ["null", "incomplete"],
     // Truthy, but not the non-empty string the schema requires.
-    ['{ supervisor_schema_version: true, dags: { d: { tasks: ["a"] } } }', 
"incomplete"],
-    ['{ supervisor_schema_version: 20260616, dags: { d: { tasks: ["a"] } } }', 
"incomplete"],
+    ['{ supervisor_schema_version: true, task_handlers: { d: { tasks: ["a"] } 
} }', "incomplete"],
+    [
+      '{ supervisor_schema_version: 20260616, task_handlers: { d: { tasks: 
["a"] } } }',
+      "incomplete",
+    ],
   ])("rejects the metadata line %s", async (manifest, message) => {
     outdir = mkdtempSync(path.join(tmpdir(), "ts-pack-"));
     const entry = path.join(outdir, "malformed-entry.ts");
@@ -453,7 +552,7 @@ describe("runPack", () => {
 
     expect(stderr()).toContain('warning: dag "empty_dag" has no tasks\n');
     expect(JSON.parse(readEmbeddedMetadata(path.join(outdir, 
"bundle.min.mjs")))).toHaveProperty(
-      "dags.empty_dag.tasks",
+      "task_handlers.empty_dag.tasks",
       [],
     );
   });
@@ -476,7 +575,7 @@ describe("runPack", () => {
     await runPack([entry, "--outdir", outdir]);
 
     const metadata = JSON.parse(readEmbeddedMetadata(path.join(outdir, 
"bundle.min.mjs")));
-    expect(metadata).toHaveProperty("dags.sales_dag");
-    expect(metadata).not.toHaveProperty("dags.billing_dag");
+    expect(metadata).toHaveProperty("task_handlers.sales_dag");
+    expect(metadata).not.toHaveProperty("task_handlers.billing_dag");
   });
 });
diff --git a/ts-sdk/tests/cli/validate.test.ts 
b/ts-sdk/tests/cli/validate.test.ts
index f97dba81d4c..143a2a93d5b 100644
--- a/ts-sdk/tests/cli/validate.test.ts
+++ b/ts-sdk/tests/cli/validate.test.ts
@@ -26,9 +26,9 @@ import type { BundleManifest } from 
"../../src/coordinator/manifest.js";
 // but two UTF-16 units, so it separates code-point from .length counting.
 const ASTRAL_LETTER = "𠀀";
 
-function warningsFor(dags: BundleManifest["dags"]): string[] {
+function warningsFor(taskHandlers: BundleManifest["task_handlers"]): string[] {
   const warnings: string[] = [];
-  warnOnSuspiciousIds(dags, (message) => warnings.push(message));
+  warnOnSuspiciousIds(taskHandlers, (message) => warnings.push(message));
   return warnings;
 }
 
diff --git a/ts-sdk/tests/coordinator/runtime-manifest.test.ts 
b/ts-sdk/tests/coordinator/runtime-manifest.test.ts
index 09f9ee1a274..b7e030f43e7 100644
--- a/ts-sdk/tests/coordinator/runtime-manifest.test.ts
+++ b/ts-sdk/tests/coordinator/runtime-manifest.test.ts
@@ -38,7 +38,7 @@ describe("buildBundleManifest", () => {
     const bundle = new Bundle(buildDag("dag_a", "t1", "t3"), buildDag("dag_b", 
"t2"));
     expect(buildBundleManifest(bundle)).toEqual({
       supervisor_schema_version: SUPERVISOR_API_VERSION,
-      dags: {
+      task_handlers: {
         dag_a: { tasks: ["t1", "t3"] },
         dag_b: { tasks: ["t2"] },
       },
@@ -46,14 +46,14 @@ describe("buildBundleManifest", () => {
   });
 
   it("keeps a registered Dag without tasks visible in the manifest", () => {
-    expect(buildBundleManifest(new 
Bundle(buildDag("empty_dag"))).dags).toEqual({
+    expect(buildBundleManifest(new 
Bundle(buildDag("empty_dag"))).task_handlers).toEqual({
       empty_dag: { tasks: [] },
     });
   });
 
   it("keeps a Dag named __proto__ visible in serialized metadata", () => {
     const manifest = buildBundleManifest(new Bundle(buildDag("__proto__", 
"task")));
-    const serializedDags = JSON.parse(JSON.stringify(manifest)).dags;
+    const serializedDags = JSON.parse(JSON.stringify(manifest)).task_handlers;
 
     expect(Object.keys(serializedDags)).toEqual(["__proto__"]);
     expect(serializedDags["__proto__"]).toEqual({ tasks: ["task"] });
@@ -62,7 +62,7 @@ describe("buildBundleManifest", () => {
   it("reports only the Dags the bundle was given", () => {
     const bundle = new Bundle(buildDag("dag_a", "t1"));
     buildDag("dag_b", "t2");
-    expect(Object.keys(buildBundleManifest(bundle).dags)).toEqual(["dag_a"]);
+    
expect(Object.keys(buildBundleManifest(bundle).task_handlers)).toEqual(["dag_a"]);
   });
 
   // The server would reject these ids. The manifest keeps them and
@@ -71,7 +71,7 @@ describe("buildBundleManifest", () => {
     "keeps a dagId the server would reject visible in the manifest: %j",
     (dagId) => {
       const manifest = buildBundleManifest(new Bundle(buildDag(dagId, "t1")));
-      expect(manifest.dags[dagId]).toEqual({ tasks: ["t1"] });
+      expect(manifest.task_handlers[dagId]).toEqual({ tasks: ["t1"] });
     },
   );
 
@@ -81,7 +81,7 @@ describe("buildBundleManifest", () => {
     "keeps a taskId the server would reject visible in the manifest: %j",
     (taskId) => {
       const manifest = buildBundleManifest(new Bundle(buildDag("example_dag", 
taskId)));
-      expect(manifest.dags["example_dag"]).toEqual({ tasks: [taskId] });
+      expect(manifest.task_handlers["example_dag"]).toEqual({ tasks: [taskId] 
});
     },
   );
 
@@ -109,6 +109,6 @@ describe("startCoordinator --airflow-metadata", () => {
     expect(written.startsWith(AIRFLOW_METADATA_SENTINEL)).toBe(true);
     const payload = 
JSON.parse(written.slice(AIRFLOW_METADATA_SENTINEL.length));
     expect(payload.supervisor_schema_version).toBe(SUPERVISOR_API_VERSION);
-    expect(payload.dags).toEqual({ metadata_dag: { tasks: ["only"] } });
+    expect(payload.task_handlers).toEqual({ metadata_dag: { tasks: ["only"] } 
});
   });
 });

Reply via email to