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 2f89e806007aaa2b9143a366e33ff6c721aaf96d Author: ZHE YOU LIU <[email protected]> AuthorDate: Tue Oct 6 05:50:35 2026 +0000 Java SDK: Bound the describe run and keep its payload to the bundle JAR Every `Jar` task gained a dependency on `packDagSources`, so a sources or javadoc JAR carried the Dag sources and claimed the `Airflow-Java-SDK-Sources` attribute that the Python side picks a bundle JAR out of a directory by. Only `jar` and `shadowJar` carry them now. The describe run had no time bound, so a `main` that does not route through `Server.create(args)` and starts serving instead held the build open with no way out. It now runs under a timeout and falls back to packing the entrypoint's source alone, which is what every other describe failure already does. Say in `serve` and `serveAsync` that a `--describe-sources` server returns without connecting, and record in `sources.json` that its paths are relative to the source directory, which is also where they sit under `META-INF/airflow/sources/`. --- .../apache/airflow/sdk/plugin/AirflowSdkPlugin.kt | 10 +++- .../apache/airflow/sdk/plugin/PackDagSources.kt | 62 +++++++++++++++------- .../airflow/sdk/plugin/AirflowSdkPluginTest.kt | 23 ++++++++ .../main/kotlin/org/apache/airflow/sdk/Server.kt | 6 +++ 4 files changed, 81 insertions(+), 20 deletions(-) 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 2ecb0acb733..c3bdbc4f8fa 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 @@ -65,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") + /** * Gradle plugin for building Apache Airflow Java SDK bundles. * @@ -108,6 +111,7 @@ abstract class AirflowBundleExtension { * the `airflow-sdk` JAR instead. The bundle JAR still contains `Main-Class`. * */ + class AirflowSdkPlugin : Plugin<Project> { override fun apply(project: Project) { project.plugins.apply("java") @@ -143,7 +147,11 @@ class AirflowSdkPlugin : Plugin<Project> { task.manifest.attributes(mapOf("Main-Class" to className)) } } - if (ext.mainClass.isPresent) { + // 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. + if (ext.mainClass.isPresent && task.name in BUNDLE_JAR_TASKS) { task.dependsOn(packTask) task.from(packTask.flatMap { it.sourcesDir }) task.manifest.attributes(mapOf(SOURCES_MANIFEST_ATTRIBUTE to SOURCES_JSON_PATH)) 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 index ac271f495ad..3ba8e98fc15 100644 --- 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 @@ -38,15 +38,22 @@ 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 +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. @@ -70,6 +77,14 @@ internal class SourceLocator( 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>, @@ -90,9 +105,6 @@ internal fun sourcesJson( */ @CacheableTask abstract class PackDagSources : DefaultTask() { - @get:Inject - abstract val execOperations: ExecOperations - @get:Input @get:Optional abstract val mainClass: Property<String> @@ -161,20 +173,32 @@ abstract class PackDagSources : DefaultTask() { main: String, describe: File, ): Map<String, String> { - val output = ByteArrayOutputStream() + var output = "" 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 + 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 } @@ -186,7 +210,7 @@ abstract class PackDagSources : DefaultTask() { } describeFailed = true val why = failure ?: "it wrote no valid --describe-sources file; does main call Server.serve?" - val log = output.toString().trim() + val log = output.trim() logger.warn( "Could not read each Dag's source from '{}' ({}); only its entrypoint source is packed.{}", main, 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 index 91fb5c5d5ed..05ed4128875 100644 --- 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 @@ -192,6 +192,29 @@ class AirflowSdkPluginTest { ) } + @Test + fun leavesTheSourcesPayloadOutOfANonBundleJar( + @TempDir dir: File, + ) { + project(dir) + dir.write( + "build.gradle", + File(dir, "build.gradle").readText() + "\njava { withSourcesJar() }\n", + ) + + gradle(dir, "jar", "sourcesJar") + + val sources = File(dir, "build/libs/bundle-test-sources.jar") + assertFalse(entries(sources).any { it.startsWith("META-INF/airflow") }) + JarFile(sources).use { assertNull(it.manifest.mainAttributes.getValue("Airflow-Java-SDK-Sources")) } + JarFile(File(dir, "build/libs/bundle-test.jar")).use { + assertEquals( + "META-INF/airflow/sources.json", + it.manifest.mainAttributes.getValue("Airflow-Java-SDK-Sources"), + ) + } + } + @Test fun skipsEverythingWithoutMainClass( @TempDir dir: File, diff --git a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Server.kt b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Server.kt index c0caccf5d47..2d8447af324 100644 --- a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Server.kt +++ b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Server.kt @@ -139,6 +139,9 @@ class Server private constructor( * The call returns when the coordinator closes the connection (normally after * one task-instance execution). * + * A [Server] built from `--describe-sources` writes that file and returns + * without connecting anywhere. + * * @param bundle Bundle containing all Dags this process can execute. * * @see [serveAsync] @@ -156,6 +159,9 @@ class Server private constructor( * coordinator closes the connection (normally after one task-instance * execution). The coroutine returns once both channels have been closed. * + * A [Server] built from `--describe-sources` writes that file and returns + * without connecting anywhere. + * * Use this variant when calling from an existing coroutine scope; use the * blocking [serve] from a plain `main` method. *
