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


##########
java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/AirflowSdkPlugin.kt:
##########
@@ -107,21 +120,45 @@ 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))
           }
         }
+        // The sources payload belongs only in the JARs a bundle is assembled 
from,
+        // so a sources or javadoc JAR does not carry it.
+        if (ext.mainClass.isPresent && task.name in BUNDLE_JAR_TASKS) {
+          task.dependsOn(packTask)

Review Comment:
   Applied in 62dd96e93bb, together with the wiring change above.



##########
java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/PackDagSources.kt:
##########
@@ -0,0 +1,240 @@
+/*
+ * 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, or
+ * else the one source file of that name when the layout does not match the
+ * package.
+ */
+internal class SourceLocator(
+  private val classesDirs: Iterable<File>,
+  private val sourceDirs: Iterable<File>,
+) {
+  /** Every source file under the source dirs by file name, walked once and 
only if a lookup misses. */
+  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('/') }
+  }
+
+  /** Path of the source file relative to its source directory, or `null` if 
it cannot be found. */
+  fun locate(className: String): String? {
+    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 } ?: return null
+    val packagePath = className.substringBeforeLast('.', "").replace('.', '/')
+    val byPackage = listOfNotNull(packagePath.ifEmpty { null }, 
sourceFile).joinToString("/")
+    if (sourceDirs.any { File(it, byPackage).isFile }) return byPackage
+    // Neither javac nor Kotlin requires the package to match the directory, 
so fall back to the
+    // file name when exactly one source tree holds it. An ambiguous name is 
left unresolved.
+    return pathsByName[sourceFile]?.distinct()?.singleOrNull()
+  }
+
+  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 }
+    // A failed run leaves describeFile missing, which Gradle would otherwise 
read as unchanged.
+    outputs.upToDateWhen { describeFile.get().asFile.isFile }
+  }
+
+  @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)
+    if (entrypoint == null) {
+      logger.warn("No source file found for entrypoint class {}; the Code view 
will show no source for this JAR", main)
+    }
+    val dagPaths = linkedMapOf<String, String>()
+    declaringClasses.forEach { (dagId, className) ->
+      val path = locator.locate(className)
+      if (path == null) {
+        logger.warn("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())
+      }
+    }

Review Comment:
   Done in b18c0aa45c9: the describe file is now the result, and a non-zero 
exit or a timeout only logs a warning when the file is there and parses.



##########
java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/AirflowSdkPlugin.kt:
##########
@@ -93,12 +98,20 @@ 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
  * the `airflow-sdk` JAR instead. The bundle JAR still contains `Main-Class`.
  *
  */
+

Review Comment:
   Applied in cb3dc2d48f2.



##########
java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/AirflowSdkPlugin.kt:
##########
@@ -63,6 +65,9 @@ abstract class AirflowBundleExtension {
   abstract val fatJar: Property<Boolean>
 }
 
+/** The JAR tasks a deployable bundle is assembled from, fat or thin. */
+private val BUNDLE_JAR_TASKS = setOf("jar", "shadowJar")

Review Comment:
   Done in 62dd96e93bb: the payload is wired onto the one JAR the bundle is 
assembled from, inside the same fatJar branch that picks it.



##########
java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/AirflowSdkPlugin.kt:
##########
@@ -107,21 +120,45 @@ 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)

Review Comment:
   Applied in ba08f8e4fab.



##########
java-sdk/plugin/src/main/kotlin/org/apache/airflow/sdk/plugin/PackDagSources.kt:
##########
@@ -0,0 +1,240 @@
+/*
+ * 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

Review Comment:
   Done in a073f81ca68: airflowBundle takes a describeTimeout, and the default 
is 30 seconds.



##########
java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Server.kt:
##########
@@ -83,10 +91,13 @@ class ApiError(
  * The process exits when the coordinator closes the connection (normally after
  * one task-instance execution).
  */
-class Server(
-  private val comm: InetSocketAddress,
-  private val logs: InetSocketAddress,
+class Server private constructor(
+  private val comm: InetSocketAddress?,
+  private val logs: InetSocketAddress?,
+  private val describeSources: File?,
 ) {

Review Comment:
   Done in 88934351785: Server is a sealed class, and create returns either 
SourceDescriber or CoordinatorServer, so the nullable addresses and the checks 
around them are gone.



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