Copilot commented on code in PR #73723:
URL: https://github.com/apache/airflow/pull/73723#discussion_r4105291977
##########
ts-sdk/src/cli/bundle-encoder.ts:
##########
@@ -139,26 +180,47 @@ function encodeHeader(regions: { metadata: Buffer;
source: Buffer; executable: B
}
/**
- * Wrap the entrypoint as written in a block comment.
+ * Wrap each source in its own block comment, one per author-owned Dag file.
*
* 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.
+ *
+ * The returned regions are keyed by path and hold byte offsets within the
+ * concatenated sources buffer (not the final bundle); the header adds the
+ * bundle-relative base offset when it renders.
*/
-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) {
+function encodeSources(sources: Record<string, string>): EncodedSources {
+ const chunks: Buffer[] = [];
+ const regions: SourceRegion[] = [];
+ let offset = 0;
+
+ for (const [path, content] of Object.entries(sources)) {
+ const openMarker =
`${EMBEDDED_SOURCE_OPEN_PREFIX}${path}${EMBEDDED_SOURCE_OPEN_SUFFIX}`;
Review Comment:
The source path is interpolated directly into a block-comment opener, unlike
the source payload below. A valid path such as `dir*/reports.ts` therefore
emits `/*# airflowSource:dir*/reports.ts`, closes the comment early, and makes
the declared source/code offsets invalid. Escape the path's block-comment
terminators (or reject such paths) before building the marker.
##########
ts-sdk/src/cli/bundle-encoder.ts:
##########
@@ -42,27 +44,45 @@ 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 EMBEDDED_SOURCES_MAX_BYTES = 4 * 1024 * 1024;
Review Comment:
The encoder now permits up to 4 MiB of source regions, while the current
task-sdk reader still rejects any source payload over 1 MiB
(`task-sdk/src/airflow/sdk/coordinators/node/_bundle_reader.py:50,430-432`).
After the framing change is made readable, a bundle between these limits will
pack successfully and then be rejected by `NodeCoordinator`; keep the producer
and consumer limits aligned or version the contract together.
This issue also appears on line 102 of the same file.
##########
ts-sdk/src/cli/pack.ts:
##########
@@ -215,8 +300,15 @@ export async function runPack(argv: readonly string[]):
Promise<void> {
platform: "node",
format: "esm",
target: "node22",
- // A digest is only worth taking over an artifact nobody reads or edits
in place.
- minify: true,
+ // Whitespace and syntax only. Identifiers stay: `dag.task(handler)`
takes
+ // the task id from the handler's name at run time, and `keepNames`
carries
+ // that name through any bundler-collision renames.
Review Comment:
This rationale is inconsistent with the SDK API: `Dag.task` receives an
explicit task ID, and `TaskHandler` documentation states that function names
are not used for task identity. Please correct the explanation (or restore full
minification if `keepNames` was added only for this assumption), because the
current comment misleadingly justifies a bundle-size change.
##########
ts-sdk/src/cli/pack.ts:
##########
@@ -202,11 +216,82 @@ async function loadEsbuild(): Promise<typeof
import("esbuild")> {
}
}
+/** Author source files esbuild's onLoad tags with their path. Same filter set
the previous
+ * task-id plugin used, which matched every language a Dag can be declared
in. */
+const AUTHOR_SOURCE_FILTER = /\.[cm]?[jt]sx?$/;
+
+/** Skip anything under node_modules: an installed dependency is not a Dag
file, and tagging it
+ * would only cost bundle bytes while overwriting the slot with the wrong
path. */
+const DEPENDENCY_PATH = /[\\/]node_modules[\\/]/;
+
+/**
+ * esbuild plugin: prepend a single line to each author-owned source file that
+ * writes the file's path into the SDK's module-source slot. `Dag`'s
+ * constructor reads that slot, so a Dag declared in `src/dags/reports.ts`
+ * carries that path even though esbuild will soon inline every module into
+ * one bundle.
+ *
+ * The prepend is text-only — no parsing, no AST — so it can never mis-identify
+ * a construction site. ES modules hoist imports above non-import statements
+ * at execution, so the tag runs *after* imported modules have written their
+ * own slots and *before* this module's own top-level statements, which is
+ * exactly when its `new Dag(...)` calls fire.
+ */
+function moduleSourceTagPlugin(cwd: string): import("esbuild").Plugin {
+ const slotKey = JSON.stringify(MODULE_SOURCE_SLOT_KEY);
+ return {
+ name: "airflow-module-source-tag",
+ setup(build) {
+ build.onLoad({ filter: AUTHOR_SOURCE_FILTER }, async ({ path: file,
namespace }) => {
+ // Only the `file` namespace has a path on disk to attribute a Dag to;
+ // a virtual module from another plugin has nothing to tag.
+ if (namespace !== "file" || DEPENDENCY_PATH.test(file)) return
undefined;
+ const source = await readFileAsync(file, "utf-8");
+ // Project-relative so the bundle is portable and readable (no host
+ // filesystem prefix), matching how esbuild's own metafile keys
sources.
+ const relative = path.relative(cwd, file);
+ const tag =
`globalThis[Symbol.for(${slotKey})]=${JSON.stringify(relative)};\n`;
Review Comment:
The slot is a mutable process-global value, so it is not scoped to this
module's evaluation. If this module reaches a top-level `await`, or another
module evaluates while it is suspended, that module can overwrite the slot
before a later `new Dag(...)` runs; the resulting `dag_source_paths` then
points to an unrelated file. The current integration test covers only
synchronous constructors, so this needs a module-scoped capture/transform or an
explicit restriction on constructors after suspension.
##########
ts-sdk/src/cli/bundle-encoder.ts:
##########
@@ -72,65 +92,86 @@ interface VerifiedByteRange {
start: string;
}
+interface VerifiedSourceRegion extends VerifiedByteRange {
+ path: string;
+}
+
interface BundleHeader {
code: VerifiedByteRange;
metadata: VerifiedByteRange;
- source: VerifiedByteRange;
+ sources: VerifiedSourceRegion[];
+}
+
+/** Per-source-region metadata within the concatenated sources buffer. */
+interface SourceRegion {
+ path: string;
+ payloadStart: number;
+ payloadEnd: number;
+ sha256: string;
+}
+
+interface EncodedSources {
+ buffer: Buffer;
+ regions: SourceRegion[];
}
export function encodeBundle(input: BundleEncoderInput): Buffer {
const metadata = encodeMetadata(input);
- const source = encodeSource(input.entrypointSource);
+ const sources = encodeSources(input.entrypointSources);
const executable = encodeExecutable(input.executable);
- const header = encodeHeader({ metadata, source, executable });
+ const header = encodeHeader({ metadata, sources, executable });
- return Buffer.concat([header, metadata, source, executable]);
+ return Buffer.concat([header, metadata, sources.buffer, executable]);
}
-function encodeHeader(regions: { metadata: Buffer; source: Buffer; executable:
Buffer }): Buffer {
+function encodeHeader(regions: {
+ metadata: Buffer;
+ sources: EncodedSources;
+ 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 metadataDigest = computeSha256(metadataPayload);
+ const codeDigest = computeSha256(regions.executable);
const zeroOffset = "0".repeat(OFFSET_HEX_WIDTH);
+ // Placeholder header with zeroed offsets and real digests, to measure its
+ // length without recursing. Every source path is present, so the array
+ // length is what it will be in the final header.
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 },
+ code: { start: zeroOffset, end: zeroOffset, sha256: codeDigest },
+ metadata: { start: zeroOffset, end: zeroOffset, sha256: metadataDigest },
+ sources: regions.sources.regions.map((region) => ({
+ path: region.path,
+ start: zeroOffset,
+ end: zeroOffset,
+ sha256: region.sha256,
+ })),
});
const metadataStart = placeholderHeader.length +
Buffer.byteLength(EMBEDDED_METADATA_PREFIX);
const metadataEnd = metadataStart + metadataPayload.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 sourcesBaseOffset = placeholderHeader.length + regions.metadata.length;
+ const codeStart = sourcesBaseOffset + regions.sources.buffer.length;
const codeEnd = codeStart + regions.executable.length;
const header = renderHeader({
code: {
start: formatOffset(codeStart),
end: formatOffset(codeEnd),
- sha256: digests.code,
+ sha256: codeDigest,
},
metadata: {
start: formatOffset(metadataStart),
end: formatOffset(metadataEnd),
- sha256: digests.metadata,
- },
- source: {
- start: formatOffset(sourceStart),
- end: formatOffset(sourceEnd),
- sha256: digests.source,
+ sha256: metadataDigest,
},
+ sources: regions.sources.regions.map((region) => ({
+ path: region.path,
+ start: formatOffset(sourcesBaseOffset + region.payloadStart),
+ end: formatOffset(sourcesBaseOffset + region.payloadEnd),
+ sha256: region.sha256,
+ })),
Review Comment:
Source paths are now serialized into the layout header, but `renderHeader`
still creates that header with `Buffer.from(..., "ascii")`. A valid non-ASCII
filename such as `données.ts` is therefore silently corrupted in the header, so
the metadata-to-region path mapping no longer round-trips. Encode the header as
UTF-8 or escape paths to ASCII before rendering it.
##########
ts-sdk/src/cli/bundle-encoder.ts:
##########
@@ -42,27 +44,45 @@ 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 EMBEDDED_SOURCES_MAX_BYTES = 4 * 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";
+/**
+ * Each source is wrapped in a block comment. The path follows the marker so
+ * a reader can tell one region from another without cross-referencing the
+ * header — the header's byte ranges are the source of truth, but a human
+ * scanning the bundle can still find their file by name.
+ */
+export const EMBEDDED_SOURCE_OPEN_PREFIX = "/*# airflowSource:";
+export const EMBEDDED_SOURCE_OPEN_SUFFIX = "\n";
export const EMBEDDED_SOURCE_CLOSE = "\n#*/\n";
export interface BundleEncoderInput {
bundleManifest: BundleManifest;
sdkVersion: string;
+ /** The entrypoint the packer was invoked on; recorded in metadata so a
reader can
+ * point at the primary source file when there is no per-Dag source path to
use. */
entrypointName: string;
- entrypointSource: string;
+ /**
+ * Every author-owned source file that declares at least one Dag, keyed by
+ * the path used in `BundleManifest.dag_source_paths`. Files that only
+ * import from these (utilities, types) are not embedded — the Code tab
+ * reads what defines each Dag, not what it depends on.
+ *
+ * A bundle whose Dags are all task handlers (Python owns them) may have
+ * an empty map here.
+ */
+ entrypointSources: Record<string, string>;
executable: Uint8Array;
}
interface BundleMetadata {
airflow_bundle_metadata_version: string;
sdk: { language: string; version: string; supervisor_schema_version: string
};
- source: string;
+ entrypoint: string;
+ dag_source_paths: BundleManifest["dag_source_paths"];
Review Comment:
The bundle contract is changed incompatibly (`source` becomes `entrypoint`,
and one source section becomes `sources`) while
`AIRFLOW_BUNDLE_METADATA_VERSION` remains `"1.0"`. Once the reader is updated,
it will either reject all existing v1 bundles or need to branch on a version
that is currently indistinguishable; bump the bundle metadata version or
implement explicit backward compatibility.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]