jason810496 commented on code in PR #74096:
URL: https://github.com/apache/airflow/pull/74096#discussion_r4209908293


##########
java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/PackDagSources.kt:
##########
@@ -0,0 +1,222 @@
+/*
+ * 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.
+ */
+
+package org.apache.airflow.sdk.plugin
+
+import groovy.json.JsonOutput
+import groovy.json.JsonSlurper
+import org.gradle.api.DefaultTask
+import org.gradle.api.file.ConfigurableFileCollection
+import org.gradle.api.file.DirectoryProperty
+import org.gradle.api.file.RegularFileProperty
+import org.gradle.api.provider.Property
+import org.gradle.api.tasks.CacheableTask
+import org.gradle.api.tasks.Classpath
+import org.gradle.api.tasks.Input
+import org.gradle.api.tasks.InputFiles
+import org.gradle.api.tasks.Nested
+import org.gradle.api.tasks.Optional
+import org.gradle.api.tasks.OutputDirectory
+import org.gradle.api.tasks.OutputFile
+import org.gradle.api.tasks.PathSensitive
+import org.gradle.api.tasks.PathSensitivity
+import org.gradle.api.tasks.TaskAction
+import org.gradle.jvm.toolchain.JavaLauncher
+import java.io.File
+import java.util.concurrent.TimeUnit
+
+internal const val SOURCES_MANIFEST_ATTRIBUTE = "Airflow-Java-SDK-Sources"
+internal const val SOURCES_JSON_PATH = "META-INF/airflow/sources.json"
+internal const val SOURCES_DIR_PATH = "META-INF/airflow/sources"
+
+/**
+ * How long the describe run gets before the build gives up on it and packs
+ * the entrypoint's source alone. A `main` that does not reach
+ * `Server.create(args)` never answers `--describe-sources`, and one that
+ * starts serving instead would otherwise hold the build open.
+ */
+private const val DESCRIBE_TIMEOUT_SECONDS = 120L
+private const val DRAIN_TIMEOUT_MILLIS = 2_000L
+
+/**
+ * Finds the source file of a class by its `SourceFile` attribute: the class
+ * file's package path plus that name, looked up in each source directory.
+ */
+internal class SourceLocator(
+  private val classesDirs: Iterable<File>,
+  private val sourceDirs: Iterable<File>,
+) {
+  /** Path of the source file relative to its source directory, or `null` if 
it cannot be found. */
+  fun locate(className: String): String? {
+    val packagePath = className.substringBeforeLast('.', "").replace('.', '/')
+    val classFile =
+      classesDirs
+        .map { File(it, className.replace('.', '/') + ".class") }
+        .firstOrNull { it.isFile } ?: return null
+    val sourceFile = readSourceFile(classFile.readBytes())?.takeIf { 
it.isNotEmpty() && '/' !in it && '\\' !in it }
+    val relative = listOfNotNull(packagePath.ifEmpty { null }, sourceFile ?: 
return null).joinToString("/")
+    return relative.takeIf { path -> sourceDirs.any { File(it, path).isFile } }
+  }
+
+  fun file(relativePath: String): File = sourceDirs.map { File(it, 
relativePath) }.first { it.isFile }
+}
+
+/**
+ * The `META-INF/airflow/sources.json` body: `entrypoint_path` and each Dag ID
+ * in `dag_source_paths`.
+ *
+ * Every path is relative to the source directory the file was found in, which
+ * is also where it sits under `META-INF/airflow/sources/` in the JAR, so a
+ * reader resolves one by prefixing that directory.
+ */
+internal fun sourcesJson(
+  entrypointPath: String?,
+  dagSourcePaths: Map<String, String>,
+): String =
+  JsonOutput.prettyPrint(
+    JsonOutput.toJson(
+      linkedMapOf<String, Any>().apply {
+        entrypointPath?.let { put("entrypoint_path", it) }
+        put("dag_source_paths", dagSourcePaths)
+      },
+    ),
+  )
+
+/**
+ * Runs the bundle's `mainClass` in describe mode to learn which class declared
+ * each Dag, then collects those classes' source files, and the entry point's,
+ * for the bundle JAR to carry.
+ */
+@CacheableTask
+abstract class PackDagSources : DefaultTask() {
+  @get:Input
+  @get:Optional
+  abstract val mainClass: Property<String>
+
+  @get:Classpath
+  abstract val classesDirs: ConfigurableFileCollection
+
+  @get:Classpath
+  abstract val runtimeClasspath: ConfigurableFileCollection
+
+  @get:InputFiles
+  @get:PathSensitive(PathSensitivity.RELATIVE)
+  abstract val sourceDirs: ConfigurableFileCollection
+
+  @get:Nested
+  @get:Optional
+  abstract val launcher: Property<JavaLauncher>
+
+  @get:OutputFile
+  abstract val describeFile: RegularFileProperty
+
+  @get:OutputDirectory
+  abstract val sourcesDir: DirectoryProperty
+
+  private var describeFailed = false
+
+  init {
+    outputs.doNotCacheIf("the Dag describe run failed") { describeFailed }

Review Comment:
   Confirmed, and fixed the way you suggested:
   
   ```kotlin
   // A failed run leaves describeFile missing, which Gradle would otherwise 
read as unchanged.
   outputs.upToDateWhen { describeFile.get().asFile.isFile }
   ```
   
   A test builds twice with a `main` that writes nothing and asserts the second 
run is not `UP-TO-DATE` and warns again.
   
   Fixed in 0cb996db1b6.
   



##########
java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/PackDagSources.kt:
##########
@@ -0,0 +1,222 @@
+/*
+ * 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.
+ */
+
+package org.apache.airflow.sdk.plugin
+
+import groovy.json.JsonOutput
+import groovy.json.JsonSlurper
+import org.gradle.api.DefaultTask
+import org.gradle.api.file.ConfigurableFileCollection
+import org.gradle.api.file.DirectoryProperty
+import org.gradle.api.file.RegularFileProperty
+import org.gradle.api.provider.Property
+import org.gradle.api.tasks.CacheableTask
+import org.gradle.api.tasks.Classpath
+import org.gradle.api.tasks.Input
+import org.gradle.api.tasks.InputFiles
+import org.gradle.api.tasks.Nested
+import org.gradle.api.tasks.Optional
+import org.gradle.api.tasks.OutputDirectory
+import org.gradle.api.tasks.OutputFile
+import org.gradle.api.tasks.PathSensitive
+import org.gradle.api.tasks.PathSensitivity
+import org.gradle.api.tasks.TaskAction
+import org.gradle.jvm.toolchain.JavaLauncher
+import java.io.File
+import java.util.concurrent.TimeUnit
+
+internal const val SOURCES_MANIFEST_ATTRIBUTE = "Airflow-Java-SDK-Sources"
+internal const val SOURCES_JSON_PATH = "META-INF/airflow/sources.json"
+internal const val SOURCES_DIR_PATH = "META-INF/airflow/sources"
+
+/**
+ * How long the describe run gets before the build gives up on it and packs
+ * the entrypoint's source alone. A `main` that does not reach
+ * `Server.create(args)` never answers `--describe-sources`, and one that
+ * starts serving instead would otherwise hold the build open.
+ */
+private const val DESCRIBE_TIMEOUT_SECONDS = 120L
+private const val DRAIN_TIMEOUT_MILLIS = 2_000L
+
+/**
+ * Finds the source file of a class by its `SourceFile` attribute: the class
+ * file's package path plus that name, looked up in each source directory.
+ */
+internal class SourceLocator(
+  private val classesDirs: Iterable<File>,
+  private val sourceDirs: Iterable<File>,
+) {
+  /** Path of the source file relative to its source directory, or `null` if 
it cannot be found. */
+  fun locate(className: String): String? {
+    val packagePath = className.substringBeforeLast('.', "").replace('.', '/')
+    val classFile =
+      classesDirs
+        .map { File(it, className.replace('.', '/') + ".class") }
+        .firstOrNull { it.isFile } ?: return null
+    val sourceFile = readSourceFile(classFile.readBytes())?.takeIf { 
it.isNotEmpty() && '/' !in it && '\\' !in it }
+    val relative = listOfNotNull(packagePath.ifEmpty { null }, sourceFile ?: 
return null).joinToString("/")
+    return relative.takeIf { path -> sourceDirs.any { File(it, path).isFile } }

Review Comment:
   Right, and the failure mode is the bad kind: nothing logged, no 
`entrypoint_path`, and a Code view telling the user the JAR embeds no source 
for a JAR this plugin built. `locate` now falls back to a unique match on the 
`SourceFile` name, over an index built once and only when the package path 
misses:
   
   ```kotlin
   private val pathsByName: Map<String, List<String>> by lazy {
     sourceDirs
       .flatMap { dir -> dir.walkTopDown().filter(File::isFile).map { 
it.relativeTo(dir).invariantSeparatorsPath } }
       .groupBy { it.substringAfterLast('/') }
   }
   ```
   
   An ambiguous name is left unresolved rather than guessed. The entrypoint 
miss now warns, and the per-Dag miss moved from `info` to `warn`. A TestKit 
case covers `package com.example.dags;` living in 
`src/main/java/dags/Reports.java`.
   
   Fixed in 0cb996db1b6.
   



##########
java-sdk/plugin/src/test/kotlin/org/apache/airflow/sdk/plugin/AirflowSdkPluginTest.kt:
##########
@@ -0,0 +1,230 @@
+/*
+ * 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.
+ */
+
+package org.apache.airflow.sdk.plugin
+
+import groovy.json.JsonSlurper
+import org.gradle.testkit.runner.BuildResult
+import org.gradle.testkit.runner.GradleRunner
+import org.gradle.testkit.runner.TaskOutcome
+import org.junit.jupiter.api.Assertions.assertEquals
+import org.junit.jupiter.api.Assertions.assertFalse
+import org.junit.jupiter.api.Assertions.assertNull
+import org.junit.jupiter.api.Assertions.assertTrue
+import org.junit.jupiter.api.Test
+import org.junit.jupiter.api.io.TempDir
+import java.io.File
+import java.util.jar.JarFile
+
+class AirflowSdkPluginTest {
+  private fun main(describeBody: String) =
+    """
+    package com.example;
+
+    public class Main {
+      public static void main(String[] args) throws Exception {
+        if (args.length == 2 && args[0].equals("--describe-sources")) {
+          $describeBody
+        }
+      }
+    }
+    """.trimIndent()
+
+  private val describeDags =
+    """
+    java.nio.file.Files.write(java.nio.file.Paths.get(args[1]), (
+        "{\"orders\":\"com.example.Main\"," +
+        "\"reports\":\"com.example.dags.Reports\"," +
+        "\"reports_backfill\":\"com.example.dags.Reports\"," +
+        "\"ghost\":\"com.example.Ghost\"}").getBytes());
+    """.trimIndent()
+
+  private fun project(
+    dir: File,
+    mainBody: String = describeDags,
+    mainClass: String? = "com.example.Main",
+  ): File {
+    dir.write("settings.gradle", "rootProject.name = 'bundle-test'\n")
+    dir.write(
+      "build.gradle",
+      """
+      plugins { id 'org.apache.airflow.sdk' }
+      airflowBundle {
+        ${mainClass?.let { "mainClass = '$it'" } ?: ""}
+        fatJar = false

Review Comment:
   Added. It builds a stub `airflow-sdk` JAR carrying the schema-version 
attribute into a local repo inside the test project, runs `shadowJar`, and 
asserts both `META-INF/airflow/sources.json` and the `Airflow-Java-SDK-Sources` 
attribute.
   
   Fixed in 0cb996db1b6.
   



##########
java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/PackDagSources.kt:
##########
@@ -0,0 +1,222 @@
+/*
+ * 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.
+ */
+
+package org.apache.airflow.sdk.plugin
+
+import groovy.json.JsonOutput
+import groovy.json.JsonSlurper
+import org.gradle.api.DefaultTask
+import org.gradle.api.file.ConfigurableFileCollection
+import org.gradle.api.file.DirectoryProperty
+import org.gradle.api.file.RegularFileProperty
+import org.gradle.api.provider.Property
+import org.gradle.api.tasks.CacheableTask
+import org.gradle.api.tasks.Classpath
+import org.gradle.api.tasks.Input
+import org.gradle.api.tasks.InputFiles
+import org.gradle.api.tasks.Nested
+import org.gradle.api.tasks.Optional
+import org.gradle.api.tasks.OutputDirectory
+import org.gradle.api.tasks.OutputFile
+import org.gradle.api.tasks.PathSensitive
+import org.gradle.api.tasks.PathSensitivity
+import org.gradle.api.tasks.TaskAction
+import org.gradle.jvm.toolchain.JavaLauncher
+import java.io.File
+import java.util.concurrent.TimeUnit
+
+internal const val SOURCES_MANIFEST_ATTRIBUTE = "Airflow-Java-SDK-Sources"
+internal const val SOURCES_JSON_PATH = "META-INF/airflow/sources.json"
+internal const val SOURCES_DIR_PATH = "META-INF/airflow/sources"
+
+/**
+ * How long the describe run gets before the build gives up on it and packs
+ * the entrypoint's source alone. A `main` that does not reach
+ * `Server.create(args)` never answers `--describe-sources`, and one that
+ * starts serving instead would otherwise hold the build open.
+ */
+private const val DESCRIBE_TIMEOUT_SECONDS = 120L
+private const val DRAIN_TIMEOUT_MILLIS = 2_000L
+
+/**
+ * Finds the source file of a class by its `SourceFile` attribute: the class
+ * file's package path plus that name, looked up in each source directory.
+ */
+internal class SourceLocator(
+  private val classesDirs: Iterable<File>,
+  private val sourceDirs: Iterable<File>,
+) {
+  /** Path of the source file relative to its source directory, or `null` if 
it cannot be found. */
+  fun locate(className: String): String? {
+    val packagePath = className.substringBeforeLast('.', "").replace('.', '/')
+    val classFile =
+      classesDirs
+        .map { File(it, className.replace('.', '/') + ".class") }
+        .firstOrNull { it.isFile } ?: return null
+    val sourceFile = readSourceFile(classFile.readBytes())?.takeIf { 
it.isNotEmpty() && '/' !in it && '\\' !in it }
+    val relative = listOfNotNull(packagePath.ifEmpty { null }, sourceFile ?: 
return null).joinToString("/")
+    return relative.takeIf { path -> sourceDirs.any { File(it, path).isFile } }
+  }
+
+  fun file(relativePath: String): File = sourceDirs.map { File(it, 
relativePath) }.first { it.isFile }
+}
+
+/**
+ * The `META-INF/airflow/sources.json` body: `entrypoint_path` and each Dag ID
+ * in `dag_source_paths`.
+ *
+ * Every path is relative to the source directory the file was found in, which
+ * is also where it sits under `META-INF/airflow/sources/` in the JAR, so a
+ * reader resolves one by prefixing that directory.
+ */
+internal fun sourcesJson(
+  entrypointPath: String?,
+  dagSourcePaths: Map<String, String>,
+): String =
+  JsonOutput.prettyPrint(
+    JsonOutput.toJson(
+      linkedMapOf<String, Any>().apply {
+        entrypointPath?.let { put("entrypoint_path", it) }
+        put("dag_source_paths", dagSourcePaths)
+      },
+    ),
+  )
+
+/**
+ * Runs the bundle's `mainClass` in describe mode to learn which class declared
+ * each Dag, then collects those classes' source files, and the entry point's,
+ * for the bundle JAR to carry.
+ */
+@CacheableTask
+abstract class PackDagSources : DefaultTask() {
+  @get:Input
+  @get:Optional
+  abstract val mainClass: Property<String>
+
+  @get:Classpath
+  abstract val classesDirs: ConfigurableFileCollection
+
+  @get:Classpath
+  abstract val runtimeClasspath: ConfigurableFileCollection
+
+  @get:InputFiles
+  @get:PathSensitive(PathSensitivity.RELATIVE)
+  abstract val sourceDirs: ConfigurableFileCollection
+
+  @get:Nested
+  @get:Optional
+  abstract val launcher: Property<JavaLauncher>
+
+  @get:OutputFile
+  abstract val describeFile: RegularFileProperty
+
+  @get:OutputDirectory
+  abstract val sourcesDir: DirectoryProperty
+
+  private var describeFailed = false
+
+  init {
+    outputs.doNotCacheIf("the Dag describe run failed") { describeFailed }
+  }
+
+  @TaskAction
+  fun pack() {
+    val describe = describeFile.get().asFile
+    val root = sourcesDir.get().asFile
+    describe.delete()
+    root.deleteRecursively()
+    describe.parentFile.mkdirs()
+
+    val main = mainClass.get()
+    val declaringClasses = describeDags(main, describe)
+    val locator = SourceLocator(classesDirs.files, sourceDirs.files)
+
+    val entrypoint = locator.locate(main)
+    val dagPaths = linkedMapOf<String, String>()
+    declaringClasses.forEach { (dagId, className) ->
+      val path = locator.locate(className)
+      if (path == null) {
+        logger.info("No source file found for class {} of Dag '{}'; it falls 
back to the entrypoint", className, dagId)
+      } else {
+        dagPaths[dagId] = path
+      }
+    }
+
+    (listOfNotNull(entrypoint) + dagPaths.values).distinct().forEach { path ->
+      val target = File(root, "$SOURCES_DIR_PATH/$path")
+      target.parentFile.mkdirs()
+      locator.file(path).copyTo(target)
+    }
+    File(root, SOURCES_JSON_PATH).apply {
+      parentFile.mkdirs()
+      writeText(sourcesJson(entrypoint, dagPaths) + "\n")
+    }
+  }
+
+  private fun describeDags(
+    main: String,
+    describe: File,
+  ): Map<String, String> {
+    var output = ""
+    val failure =
+      try {
+        val java =
+          launcher.orNull
+            ?.executablePath
+            ?.asFile
+            ?.absolutePath ?: "java"
+        val classpath = (classesDirs + 
runtimeClasspath).joinToString(File.pathSeparator) { it.absolutePath }
+        val process =
+          ProcessBuilder(java, "-cp", classpath, main, "--describe-sources", 
describe.absolutePath)
+            .redirectErrorStream(true)
+            .start()
+        // Drained on its own thread: a main that writes more than the pipe
+        // holds would otherwise block before the timeout could fire.
+        val drain = Thread { output = 
process.inputStream.bufferedReader().readText() }
+        drain.isDaemon = true
+        drain.start()
+        val finished = process.waitFor(DESCRIBE_TIMEOUT_SECONDS, 
TimeUnit.SECONDS)
+        if (!finished) process.destroyForcibly()
+        drain.join(DRAIN_TIMEOUT_MILLIS)
+        when {
+          !finished -> "it did not finish within $DESCRIBE_TIMEOUT_SECONDS 
seconds"
+          process.exitValue() != 0 -> "exit code ${process.exitValue()}"
+          else -> null
+        }
+      } catch (e: Exception) {
+        e.message ?: e.javaClass.simpleName
+      }
+    if (failure == null && describe.isFile) {
+      @Suppress("UNCHECKED_CAST")
+      (runCatching { JsonSlurper().parse(describe) }.getOrNull() as? 
Map<String, Any?>)?.let { parsed ->
+        return parsed.mapNotNull { (id, cls) -> (cls as? String)?.let { id to 
it } }.toMap(linkedMapOf())
+      }
+    }
+    describeFailed = true
+    val why = failure ?: "it wrote no valid --describe-sources file; does main 
call Server.serve?"

Review Comment:
   Reworded to "does main pass its args to Server.create?", and the assertion 
in `AirflowSdkPluginTest` follows it.
   
   Fixed in 0cb996db1b6.
   



##########
java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/DagSource.kt:
##########
@@ -0,0 +1,66 @@
+/*
+ * 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.
+ */
+
+package org.apache.airflow.sdk.internal
+
+import org.apache.airflow.sdk.DagDef
+
+/**
+ * @suppress
+ *
+ * Tracks which class declared a [DagDef], so the bundle can ship that class's
+ * source file. Public so that processor-generated builders can call 
[declaredBy];
+ * not user-facing API.
+ */
+object DagSource {
+  private const val SDK_PACKAGE = "org.apache.airflow.sdk."
+  private val IGNORED_PREFIXES = listOf(SDK_PACKAGE, "java.", "javax.", 
"jdk.", "sun.", "kotlin.", "kotlinx.")
+
+  private val walker = 
StackWalker.getInstance(StackWalker.Option.RETAIN_CLASS_REFERENCE)
+
+  /**
+   * Names [declaring] as the class that declared [dag], replacing what was
+   * captured at construction. A generated builder calls this so the Dag points
+   * at the annotated class rather than the builder.
+   */
+  @JvmStatic
+  fun declaredBy(
+    dag: DagDef,
+    declaring: Class<*>,
+  ): DagDef {
+    dag.declaringClass = outermost(declaring)
+    return dag
+  }
+
+  /** The outermost class of the first caller outside the SDK and the standard 
libraries. */
+  internal fun capture(): Class<*>? =
+    walker.walk { frames ->
+      frames
+        .map { it.declaringClass }
+        .filter { c -> IGNORED_PREFIXES.none { c.name.startsWith(it) } }
+        .findFirst()

Review Comment:
   Intended, and now documented rather than changed. The JVM has no module body 
for ADR-0015 decision 1 to key on, so a Dag is recorded against the class that 
constructed it. `java.rst` says so, including that a factory-built Dag maps to 
the factory's class and that a factory in a dependency JAR has no source in the 
project and falls back to the `mainClass` source. ADR-0015 decision 5 notes why 
a JVM SDK resolves differently.
   
   I deliberately left `DagSource.declaredBy` out of the docs: it is reachable, 
but it sits in the `internal` package with `@suppress` on its KDoc and exists 
for processor-generated builders, so documenting it as a workaround would be 
promoting it to API. Say the word if you would rather it became supported.
   
   Documented in 0cb996db1b6.
   



##########
java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/AirflowSdkPlugin.kt:
##########
@@ -107,21 +120,47 @@ class AirflowSdkPlugin : Plugin<Project> {
     ext.fatJar.convention(true)
 
     project.afterEvaluate {
+      val main = 
project.extensions.getByType(SourceSetContainer::class.java).getByName("main")
+
+      val packTask =
+        project.tasks.register("packDagSources", PackDagSources::class.java) { 
task ->
+          task.group = "build"
+          task.description = "Collects the source file of each Dag and the 
entrypoint to pack into the bundle JAR."
+          task.dependsOn(project.tasks.named("classes"))
+          task.onlyIf { ext.mainClass.isPresent }
+          task.mainClass.set(ext.mainClass)
+          task.classesDirs.from(main.output.classesDirs)
+          task.runtimeClasspath.from(main.output, main.runtimeClasspath)
+          task.sourceDirs.from(main.allSource.srcDirs)
+          task.launcher.convention(
+            project.extensions
+              .getByType(JavaToolchainService::class.java)
+              
.launcherFor(project.extensions.getByType(JavaPluginExtension::class.java).toolchain),
+          )
+          
task.describeFile.set(project.layout.buildDirectory.file("airflow/describe-sources.json"))
+          
task.sourcesDir.set(project.layout.buildDirectory.dir("airflow/sources"))
+        }
+
       project.tasks.withType(Jar::class.java).configureEach { task ->
         task.doFirst {
           ext.mainClass.orNull?.let { className ->
             task.manifest.attributes(mapOf("Main-Class" to className))
           }
         }
+        // Only the JARs a bundle is assembled from carry the Dag sources. A
+        // sources or javadoc JAR would claim the attribute without being a
+        // bundle, and the Python side picks the bundle JAR out of a directory
+        // by that attribute.

Review Comment:
   Right, `Main-Class` is what both `might_contain_dag` and the coordinator 
select on, and the `doFirst` stamps it on every `Jar` task anyway. The comment 
now just says the payload belongs only in the JARs a bundle is assembled from.
   
   Fixed in 0cb996db1b6.
   



##########
java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/PackDagSources.kt:
##########
@@ -0,0 +1,222 @@
+/*
+ * 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.
+ */
+
+package org.apache.airflow.sdk.plugin
+
+import groovy.json.JsonOutput
+import groovy.json.JsonSlurper
+import org.gradle.api.DefaultTask
+import org.gradle.api.file.ConfigurableFileCollection
+import org.gradle.api.file.DirectoryProperty
+import org.gradle.api.file.RegularFileProperty
+import org.gradle.api.provider.Property
+import org.gradle.api.tasks.CacheableTask
+import org.gradle.api.tasks.Classpath
+import org.gradle.api.tasks.Input
+import org.gradle.api.tasks.InputFiles
+import org.gradle.api.tasks.Nested
+import org.gradle.api.tasks.Optional
+import org.gradle.api.tasks.OutputDirectory
+import org.gradle.api.tasks.OutputFile
+import org.gradle.api.tasks.PathSensitive
+import org.gradle.api.tasks.PathSensitivity
+import org.gradle.api.tasks.TaskAction
+import org.gradle.jvm.toolchain.JavaLauncher
+import java.io.File
+import java.util.concurrent.TimeUnit
+
+internal const val SOURCES_MANIFEST_ATTRIBUTE = "Airflow-Java-SDK-Sources"
+internal const val SOURCES_JSON_PATH = "META-INF/airflow/sources.json"
+internal const val SOURCES_DIR_PATH = "META-INF/airflow/sources"
+
+/**
+ * How long the describe run gets before the build gives up on it and packs
+ * the entrypoint's source alone. A `main` that does not reach
+ * `Server.create(args)` never answers `--describe-sources`, and one that
+ * starts serving instead would otherwise hold the build open.
+ */
+private const val DESCRIBE_TIMEOUT_SECONDS = 120L
+private const val DRAIN_TIMEOUT_MILLIS = 2_000L
+
+/**
+ * Finds the source file of a class by its `SourceFile` attribute: the class
+ * file's package path plus that name, looked up in each source directory.
+ */
+internal class SourceLocator(
+  private val classesDirs: Iterable<File>,
+  private val sourceDirs: Iterable<File>,
+) {
+  /** Path of the source file relative to its source directory, or `null` if 
it cannot be found. */
+  fun locate(className: String): String? {
+    val packagePath = className.substringBeforeLast('.', "").replace('.', '/')
+    val classFile =
+      classesDirs
+        .map { File(it, className.replace('.', '/') + ".class") }
+        .firstOrNull { it.isFile } ?: return null
+    val sourceFile = readSourceFile(classFile.readBytes())?.takeIf { 
it.isNotEmpty() && '/' !in it && '\\' !in it }
+    val relative = listOfNotNull(packagePath.ifEmpty { null }, sourceFile ?: 
return null).joinToString("/")
+    return relative.takeIf { path -> sourceDirs.any { File(it, path).isFile } }
+  }
+
+  fun file(relativePath: String): File = sourceDirs.map { File(it, 
relativePath) }.first { it.isFile }
+}
+
+/**
+ * The `META-INF/airflow/sources.json` body: `entrypoint_path` and each Dag ID
+ * in `dag_source_paths`.
+ *
+ * Every path is relative to the source directory the file was found in, which
+ * is also where it sits under `META-INF/airflow/sources/` in the JAR, so a
+ * reader resolves one by prefixing that directory.
+ */
+internal fun sourcesJson(
+  entrypointPath: String?,
+  dagSourcePaths: Map<String, String>,
+): String =
+  JsonOutput.prettyPrint(
+    JsonOutput.toJson(
+      linkedMapOf<String, Any>().apply {
+        entrypointPath?.let { put("entrypoint_path", it) }
+        put("dag_source_paths", dagSourcePaths)

Review Comment:
   Good catch, `dag_source_paths` was written and never read. `get_source_code` 
now passes the ID it is given:
   
   ```python
   info = _find_source_entry(zf, dag_id)
   ```
   
   Tests cover a Dag mapped to its own file getting that file back, an unmapped 
Dag, and `dag_id=None`, with the entrypoint fallback where expected.
   
   Fixed in 0cb996db1b6.
   



-- 
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]

Reply via email to