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

jason810496 pushed a commit to branch jason/java-sdk/pack-dag-sources
in repository https://gitbox.apache.org/repos/asf/airflow.git

commit b6f535918f34ebda1422375f1844f15d0965364a
Author: ZHE YOU LIU <[email protected]>
AuthorDate: Fri Oct 2 15:19:36 2026 +0000

    Java SDK: Pack each native Dag's source file into the bundle JAR
    
    The packDagSources task runs mainClass in describe mode, finds each Dag's
    source file through the class file's SourceFile attribute, and packs each
    unique file with a sources.json mapping under META-INF/airflow/. The JAR
    manifest names it with Airflow-Java-SDK-Sources.
---
 java-sdk/README.md                                 |   4 +
 java-sdk/plugin/build.gradle.kts                   |   6 +
 .../apache/airflow/sdk/plugin/AirflowSdkPlugin.kt  |  41 +++-
 .../org/apache/airflow/sdk/plugin/ClassFiles.kt    |  85 +++++++++
 .../apache/airflow/sdk/plugin/PackDagSources.kt    | 198 ++++++++++++++++++++
 .../airflow/sdk/plugin/AirflowSdkPluginTest.kt     | 207 +++++++++++++++++++++
 .../apache/airflow/sdk/plugin/ClassFilesTest.kt    |  80 ++++++++
 .../airflow/sdk/plugin/PackDagSourcesTest.kt       |  89 +++++++++
 .../org/apache/airflow/sdk/plugin/TestSupport.kt   |  45 +++++
 9 files changed, 750 insertions(+), 5 deletions(-)

diff --git a/java-sdk/README.md b/java-sdk/README.md
index 57e71ee3ebb..9e28064f5e5 100644
--- a/java-sdk/README.md
+++ b/java-sdk/README.md
@@ -91,6 +91,10 @@ Now `cd example` into the example project, and
   ../gradlew bundle
   ```
 
+  The bundle JAR also carries the source file of each Dag declared in Java
+  (and of the entrypoint) under `META-INF/airflow/`, so Airflow can show a
+  Dag's source.
+
 * Put the [DAG with stub tasks](./example/src/resources/dags) to somewhere 
Airflow can find.
 
 * Ensure the `java` command is available in the same environment the Airflow
diff --git a/java-sdk/plugin/build.gradle.kts b/java-sdk/plugin/build.gradle.kts
index edf930bb6fc..7efcdc36e61 100644
--- a/java-sdk/plugin/build.gradle.kts
+++ b/java-sdk/plugin/build.gradle.kts
@@ -25,6 +25,12 @@ plugins {
 
 dependencies {
     implementation("com.gradleup.shadow:shadow-gradle-plugin:9.1.0") // Last 
supporting Java 11.
+
+    testImplementation(kotlin("test"))
+}
+
+tasks.withType<Test> {
+    useJUnitPlatform()
 }
 
 gradlePlugin {
diff --git 
a/java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/AirflowSdkPlugin.kt
 
b/java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/AirflowSdkPlugin.kt
index 95836c12fc5..fd920b06430 100644
--- 
a/java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/AirflowSdkPlugin.kt
+++ 
b/java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/AirflowSdkPlugin.kt
@@ -22,11 +22,13 @@ package org.apache.airflow.sdk.plugin
 import com.github.jengelman.gradle.plugins.shadow.tasks.ShadowJar
 import org.gradle.api.Plugin
 import org.gradle.api.Project
+import org.gradle.api.plugins.JavaPluginExtension
 import org.gradle.api.provider.Property
 import org.gradle.api.tasks.Copy
 import org.gradle.api.tasks.Input
 import org.gradle.api.tasks.SourceSetContainer
 import org.gradle.api.tasks.bundling.Jar
+import org.gradle.jvm.toolchain.JavaToolchainService
 import java.lang.reflect.Modifier
 import java.net.URLClassLoader
 import java.util.jar.JarFile
@@ -93,6 +95,13 @@ abstract class AirflowBundleExtension {
  * identify which version of the Supervisor Schema it should use to communicate
  * with the built JAR.
  *
+ * When `mainClass` is set, the bundle JAR also carries the source file of each
+ * Dag declared in Java and of the entrypoint, so Airflow can show a Dag's 
source.
+ * They are collected by the `packDagSources` task, which runs `mainClass` 
once in
+ * a describe mode, and are listed in `META-INF/airflow/sources.json`, named 
by the
+ * `Airflow-Java-SDK-Sources` manifest attribute. If that run fails, a warning 
is
+ * logged and only the entrypoint's source is packed.
+ *
  * If `fatJar` is explicitly set to `false`, the `bundle` task builds a bare 
JAR
  * containing only the Dag bundle, and collect all dependency JARs into the 
target
  * directory instead. In this mode, `Airflow-Supervisor-Schema-Version` lives 
in
@@ -107,21 +116,43 @@ 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))
           }
         }
+        if (ext.mainClass.isPresent) {
+          task.dependsOn(packTask)
+          task.from(packTask.flatMap { it.sourcesDir })
+          task.manifest.attributes(mapOf(SOURCES_MANIFEST_ATTRIBUTE to 
SOURCES_JSON_PATH))
+        }
       }
 
       val classFiles =
         project.objects.fileCollection().from(
-          project.extensions
-            .getByType(SourceSetContainer::class.java)
-            .getByName("main")
-            .output
-            .classesDirs,
+          main.output.classesDirs,
           project.configurations.getByName("runtimeClasspath"),
         )
 
diff --git 
a/java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/ClassFiles.kt 
b/java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/ClassFiles.kt
new file mode 100644
index 00000000000..ef693559501
--- /dev/null
+++ 
b/java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/ClassFiles.kt
@@ -0,0 +1,85 @@
+/*
+ * 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 java.io.ByteArrayInputStream
+import java.io.DataInputStream
+import java.io.IOException
+
+private const val CLASS_MAGIC = 0xCAFEBABE.toInt()
+
+/**
+ * Reads the `SourceFile` attribute of a class file, which javac, kotlinc and 
scalac all
+ * write, or returns `null` if the class has none or the bytes are not a class 
file.
+ */
+internal fun readSourceFile(classFile: ByteArray): String? =
+  try {
+    DataInputStream(ByteArrayInputStream(classFile)).use(::readSourceFile)
+  } catch (_: IOException) {
+    null
+  }
+
+private fun readSourceFile(input: DataInputStream): String? {
+  if (input.readInt() != CLASS_MAGIC) return null
+  input.skipBytes(4) // minor and major version
+
+  val utf8 = arrayOfNulls<String>(input.readUnsignedShort())
+  var index = 1
+  while (index < utf8.size) {
+    when (val tag = input.readUnsignedByte()) {
+      1 -> utf8[index] = input.readUTF()
+      3, 4, 9, 10, 11, 12, 17, 18 -> input.skipBytes(4)
+      5, 6 -> {
+        input.skipBytes(8)
+        index++ // long and double take two slots
+      }
+      7, 8, 16, 19, 20 -> input.skipBytes(2)
+      15 -> input.skipBytes(3)
+      else -> throw IOException("Unknown constant pool tag $tag")
+    }
+    index++
+  }
+
+  input.skipBytes(6) // access flags, this class, super class
+  input.skipBytes(2 * input.readUnsignedShort()) // interfaces
+  repeat(2) { skipMembers(input) } // fields, then methods
+
+  repeat(input.readUnsignedShort()) {
+    val name = utf8.getOrNull(input.readUnsignedShort())
+    val length = input.readInt()
+    if (name == "SourceFile" && length == 2) return 
utf8.getOrNull(input.readUnsignedShort())
+    input.skipBytes(length)
+  }
+  return null
+}
+
+private fun skipMembers(input: DataInputStream) {
+  repeat(input.readUnsignedShort()) {
+    input.skipBytes(6) // access flags, name, descriptor
+    skipAttributes(input)
+  }
+}
+
+private fun skipAttributes(input: DataInputStream) {
+  repeat(input.readUnsignedShort()) {
+    input.skipBytes(2)
+    input.skipBytes(input.readInt())
+  }
+}
diff --git 
a/java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/PackDagSources.kt
 
b/java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/PackDagSources.kt
new file mode 100644
index 00000000000..ac271f495ad
--- /dev/null
+++ 
b/java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/PackDagSources.kt
@@ -0,0 +1,198 @@
+/*
+ * 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 org.gradle.process.ExecOperations
+import java.io.ByteArrayOutputStream
+import java.io.File
+import javax.inject.Inject
+
+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"
+
+/**
+ * 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 }
+}
+
+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:Inject
+  abstract val execOperations: ExecOperations
+
+  @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> {
+    val output = ByteArrayOutputStream()
+    val failure =
+      try {
+        val result =
+          execOperations.javaexec { spec ->
+            launcher.orNull?.let { 
spec.executable(it.executablePath.asFile.absolutePath) }
+            spec.classpath(classesDirs, runtimeClasspath)
+            spec.mainClass.set(main)
+            spec.args("--describe-sources", describe.absolutePath)
+            spec.isIgnoreExitValue = true
+            spec.standardOutput = output
+            spec.errorOutput = output
+          }
+        if (result.exitValue != 0) "exit code ${result.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?"
+    val log = output.toString().trim()
+    logger.warn(
+      "Could not read each Dag's source from '{}' ({}); only its entrypoint 
source is packed.{}",
+      main,
+      why,
+      if (log.isEmpty()) "" else "\n$log",
+    )
+    return emptyMap()
+  }
+}
diff --git 
a/java-sdk/plugin/src/test/kotlin/org/apache/airflow/sdk/plugin/AirflowSdkPluginTest.kt
 
b/java-sdk/plugin/src/test/kotlin/org/apache/airflow/sdk/plugin/AirflowSdkPluginTest.kt
new file mode 100644
index 00000000000..91fb5c5d5ed
--- /dev/null
+++ 
b/java-sdk/plugin/src/test/kotlin/org/apache/airflow/sdk/plugin/AirflowSdkPluginTest.kt
@@ -0,0 +1,207 @@
+/*
+ * 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
+      }
+      sourceSets { main { java.srcDir 'src/extra' } }
+      """.trimIndent(),
+    )
+    dir.write("src/main/java/com/example/Main.java", main(mainBody))
+    dir.write("src/extra/com/example/dags/Reports.java", "package 
com.example.dags; public class Reports {}")
+    return dir
+  }
+
+  private fun gradle(
+    dir: File,
+    vararg tasks: String,
+  ): BuildResult =
+    GradleRunner
+      .create()
+      .withPluginClasspath()
+      .withProjectDir(dir)
+      .withArguments(*tasks, "--stacktrace")
+      .build()
+
+  private fun entries(jar: File): List<String> =
+    JarFile(jar).use {
+      it
+        .entries()
+        .asSequence()
+        .map { e -> e.name }
+        .toList()
+    }
+
+  @Suppress("UNCHECKED_CAST")
+  private fun sourcesJson(jar: File): Map<String, Any> =
+    JarFile(jar).use {
+      
JsonSlurper().parse(it.getInputStream(it.getJarEntry("META-INF/airflow/sources.json")))
 as Map<String, Any>
+    }
+
+  @Test
+  fun packsEachDagSourceOnceAndTheEntrypoint(
+    @TempDir dir: File,
+  ) {
+    project(dir)
+    val result = gradle(dir, "jar")
+    assertEquals(TaskOutcome.SUCCESS, result.task(":packDagSources")!!.outcome)
+
+    val jar = File(dir, "build/libs/bundle-test.jar")
+    assertEquals(
+      listOf(
+        "META-INF/airflow/sources.json",
+        "META-INF/airflow/sources/com/example/Main.java",
+        "META-INF/airflow/sources/com/example/dags/Reports.java",
+      ),
+      entries(jar).filter { it.startsWith("META-INF/airflow/") && 
!it.endsWith("/") }.sorted(),
+    )
+    assertEquals(
+      mapOf(
+        "entrypoint_path" to "com/example/Main.java",
+        "dag_source_paths" to
+          mapOf(
+            "orders" to "com/example/Main.java",
+            "reports" to "com/example/dags/Reports.java",
+            "reports_backfill" to "com/example/dags/Reports.java",
+          ),
+      ),
+      sourcesJson(jar),
+    )
+    JarFile(jar).use {
+      assertEquals("META-INF/airflow/sources.json", 
it.manifest.mainAttributes.getValue("Airflow-Java-SDK-Sources"))
+      assertEquals("com.example.Main", 
it.manifest.mainAttributes.getValue("Main-Class"))
+    }
+  }
+
+  @Test
+  fun isUpToDateOnSecondRunAndRerunsWhenASourceChanges(
+    @TempDir dir: File,
+  ) {
+    project(dir)
+    gradle(dir, "jar")
+    assertEquals(TaskOutcome.UP_TO_DATE, gradle(dir, 
"jar").task(":packDagSources")!!.outcome)
+
+    dir.write("src/extra/com/example/dags/Reports.java", "package 
com.example.dags; public class Reports { int v; }")
+    assertEquals(TaskOutcome.SUCCESS, gradle(dir, 
"jar").task(":packDagSources")!!.outcome)
+    val packed = File(dir, 
"build/airflow/sources/META-INF/airflow/sources/com/example/dags/Reports.java")
+    assertTrue(packed.readText().contains("int v"))
+  }
+
+  @Test
+  fun worksWithConfigurationCache(
+    @TempDir dir: File,
+  ) {
+    project(dir)
+    gradle(dir, "jar", "--configuration-cache")
+    val again = gradle(dir, "jar", "--configuration-cache")
+    assertTrue(again.output.contains("Reusing configuration cache"))
+  }
+
+  @Test
+  fun fallsBackToEntrypointOnlyWhenTheDescribeRunFails(
+    @TempDir dir: File,
+  ) {
+    project(dir, mainBody = """throw new IllegalStateException("boom");""")
+    val result = gradle(dir, "jar")
+
+    assertEquals(TaskOutcome.SUCCESS, result.task(":jar")!!.outcome)
+    assertTrue(result.output.contains("only its entrypoint source is packed"), 
result.output)
+    assertTrue(result.output.contains("boom"), result.output)
+    assertEquals(
+      mapOf("entrypoint_path" to "com/example/Main.java", "dag_source_paths" 
to emptyMap<String, String>()),
+      sourcesJson(File(dir, "build/libs/bundle-test.jar")),
+    )
+  }
+
+  @Test
+  fun fallsBackToEntrypointOnlyWhenMainWritesNothing(
+    @TempDir dir: File,
+  ) {
+    project(dir, mainBody = "")
+    val result = gradle(dir, "jar")
+
+    assertTrue(result.output.contains("does main call Server.serve?"), 
result.output)
+    assertEquals(
+      mapOf("entrypoint_path" to "com/example/Main.java", "dag_source_paths" 
to emptyMap<String, String>()),
+      sourcesJson(File(dir, "build/libs/bundle-test.jar")),
+    )
+  }
+
+  @Test
+  fun skipsEverythingWithoutMainClass(
+    @TempDir dir: File,
+  ) {
+    project(dir, mainClass = null)
+    val result = gradle(dir, "jar", "packDagSources")
+
+    assertEquals(TaskOutcome.SKIPPED, result.task(":packDagSources")!!.outcome)
+    val jar = File(dir, "build/libs/bundle-test.jar")
+    assertFalse(entries(jar).any { it.startsWith("META-INF/airflow") })
+    JarFile(jar).use { 
assertNull(it.manifest.mainAttributes.getValue("Airflow-Java-SDK-Sources")) }
+  }
+}
diff --git 
a/java-sdk/plugin/src/test/kotlin/org/apache/airflow/sdk/plugin/ClassFilesTest.kt
 
b/java-sdk/plugin/src/test/kotlin/org/apache/airflow/sdk/plugin/ClassFilesTest.kt
new file mode 100644
index 00000000000..693a8a21c12
--- /dev/null
+++ 
b/java-sdk/plugin/src/test/kotlin/org/apache/airflow/sdk/plugin/ClassFilesTest.kt
@@ -0,0 +1,80 @@
+/*
+ * 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 org.junit.jupiter.api.Assertions.assertEquals
+import org.junit.jupiter.api.Assertions.assertNull
+import org.junit.jupiter.api.Test
+import org.junit.jupiter.api.io.TempDir
+import java.io.File
+
+class ClassFilesTest {
+  private class Nested
+
+  private fun bytesOf(cls: Class<*>): ByteArray =
+    cls.getResourceAsStream("/" + cls.name.replace('.', '/') + ".class")!!.use 
{ it.readBytes() }
+
+  @Test
+  fun readsSourceFileOfKotlinClass() {
+    assertEquals("ClassFilesTest.kt", 
readSourceFile(bytesOf(ClassFilesTest::class.java)))
+  }
+
+  @Test
+  fun readsSourceFileOfNestedClassFromItsOutermostFile() {
+    assertEquals("ClassFilesTest.kt", 
readSourceFile(bytesOf(Nested::class.java)))
+  }
+
+  @Test
+  fun readsSourceFileOfJdkClassesWithLongAndDoubleConstants() {
+    // These carry every constant pool kind, including long, double and 
invokedynamic entries.
+    assertEquals("String.java", readSourceFile(bytesOf(String::class.java)))
+    assertEquals("Long.java", 
readSourceFile(bytesOf(Long::class.javaObjectType)))
+    assertEquals("Double.java", 
readSourceFile(bytesOf(Double::class.javaObjectType)))
+    assertEquals("HashMap.java", 
readSourceFile(bytesOf(java.util.HashMap::class.java)))
+  }
+
+  @Test
+  fun readsSourceFileOfJavaClassWithSeveralTopLevelTypes(
+    @TempDir dir: File,
+  ) {
+    dir.write("p/Main.java", "package p; public class Main { static final long 
L = 1L; } class Other {}")
+    compileJava(dir, File(dir, "out"), "p/Main.java")
+
+    assertEquals("Main.java", readSourceFile(File(dir, 
"out/p/Main.class").readBytes()))
+    assertEquals("Main.java", readSourceFile(File(dir, 
"out/p/Other.class").readBytes()))
+  }
+
+  @Test
+  fun returnsNullWithoutSourceFileAttribute(
+    @TempDir dir: File,
+  ) {
+    dir.write("p/Main.java", "package p; public class Main {}")
+    compileJava(dir, File(dir, "out"), "p/Main.java", options = 
listOf("-g:none"))
+
+    assertNull(readSourceFile(File(dir, "out/p/Main.class").readBytes()))
+  }
+
+  @Test
+  fun returnsNullForBytesThatAreNotAClassFile() {
+    assertNull(readSourceFile(ByteArray(0)))
+    assertNull(readSourceFile("not a class file".toByteArray()))
+    assertNull(readSourceFile(bytesOf(ClassFilesTest::class.java).copyOf(40)))
+  }
+}
diff --git 
a/java-sdk/plugin/src/test/kotlin/org/apache/airflow/sdk/plugin/PackDagSourcesTest.kt
 
b/java-sdk/plugin/src/test/kotlin/org/apache/airflow/sdk/plugin/PackDagSourcesTest.kt
new file mode 100644
index 00000000000..bca13df2396
--- /dev/null
+++ 
b/java-sdk/plugin/src/test/kotlin/org/apache/airflow/sdk/plugin/PackDagSourcesTest.kt
@@ -0,0 +1,89 @@
+/*
+ * 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 org.junit.jupiter.api.Assertions.assertEquals
+import org.junit.jupiter.api.Assertions.assertNull
+import org.junit.jupiter.api.Test
+import org.junit.jupiter.api.io.TempDir
+import java.io.File
+
+class PackDagSourcesTest {
+  private fun project(dir: File): Pair<File, File> {
+    val src = File(dir, "src")
+    val extra = File(dir, "extra")
+    src.write("com/example/Main.java", "package com.example; public class Main 
{}")
+    src.write("com/example/Same.java", "package com.example; public class Same 
{ class Inner {} }")
+    extra.write("com/example/dags/Reports.java", "package com.example.dags; 
public class Reports {} class Helper {}")
+    src.write("com/example/NoSource.java", "package com.example; public class 
NoSource {}")
+    val classes = File(dir, "classes")
+    compileJava(
+      src,
+      classes,
+      "com/example/Main.java",
+      "com/example/Same.java",
+      "com/example/NoSource.java",
+    )
+    compileJava(extra, classes, "com/example/dags/Reports.java")
+    File(src, "com/example/NoSource.java").delete()
+    return src to classes
+  }
+
+  @Test
+  fun locatesSourceByPackagePathAndSourceFileAttribute(
+    @TempDir dir: File,
+  ) {
+    val (src, classes) = project(dir)
+    val locator = SourceLocator(listOf(classes), listOf(src, File(dir, 
"extra")))
+
+    assertEquals("com/example/Main.java", locator.locate("com.example.Main"))
+    assertEquals("com/example/Same.java", 
locator.locate("com.example.Same\$Inner"))
+    assertEquals("com/example/dags/Reports.java", 
locator.locate("com.example.dags.Reports"))
+    assertEquals("com/example/dags/Reports.java", 
locator.locate("com.example.dags.Helper"))
+    assertEquals(File(dir, "extra/com/example/dags/Reports.java"), 
locator.file("com/example/dags/Reports.java"))
+  }
+
+  @Test
+  fun locatesNothingForUnknownClassOrMissingSource(
+    @TempDir dir: File,
+  ) {
+    val (src, classes) = project(dir)
+    val locator = SourceLocator(listOf(classes), listOf(src))
+
+    assertNull(locator.locate("com.example.Missing"))
+    assertNull(locator.locate("com.example.NoSource"))
+    assertNull(locator.locate("com.example.dags.Reports"))
+  }
+
+  @Test
+  fun rendersSourcesJson() {
+    assertEquals(
+      
"""{"entrypoint_path":"com/example/Main.java","dag_source_paths":{"orders":"com/example/Main.java"}}""",
+      compact(sourcesJson("com/example/Main.java", mapOf("orders" to 
"com/example/Main.java"))),
+    )
+  }
+
+  @Test
+  fun omitsEntrypointPathWhenUnresolved() {
+    assertEquals("""{"dag_source_paths":{}}""", compact(sourcesJson(null, 
emptyMap())))
+  }
+
+  private fun compact(json: String) = 
(groovy.json.JsonSlurper().parseText(json) as Map<*, 
*>).let(groovy.json.JsonOutput::toJson)
+}
diff --git 
a/java-sdk/plugin/src/test/kotlin/org/apache/airflow/sdk/plugin/TestSupport.kt 
b/java-sdk/plugin/src/test/kotlin/org/apache/airflow/sdk/plugin/TestSupport.kt
new file mode 100644
index 00000000000..d8acc86fcfc
--- /dev/null
+++ 
b/java-sdk/plugin/src/test/kotlin/org/apache/airflow/sdk/plugin/TestSupport.kt
@@ -0,0 +1,45 @@
+/*
+ * 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 java.io.File
+import javax.tools.ToolProvider
+
+/** Compiles [sources] (paths under [srcDir]) into [classesDir] with the JDK's 
own compiler. */
+internal fun compileJava(
+  srcDir: File,
+  classesDir: File,
+  vararg sources: String,
+  options: List<String> = emptyList(),
+) {
+  classesDir.mkdirs()
+  val compiler = checkNotNull(ToolProvider.getSystemJavaCompiler()) { "Tests 
need a JDK" }
+  val args = options + listOf("-d", classesDir.path) + sources.map { 
File(srcDir, it).path }
+  check(compiler.run(null, null, null, *args.toTypedArray()) == 0) { "javac 
failed for ${sources.toList()}" }
+}
+
+internal fun File.write(
+  relativePath: String,
+  text: String,
+): File =
+  File(this, relativePath).apply {
+    parentFile.mkdirs()
+    writeText(text)
+  }

Reply via email to