This is an automated email from the ASF dual-hosted git repository.
Abacn pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new ef94119c2f3 Fix artifact staging windows filename (#39916)
ef94119c2f3 is described below
commit ef94119c2f37eed92cf4b8968b98bdb9c5a0e019
Author: Yi Hu <[email protected]>
AuthorDate: Tue Sep 8 14:59:40 2026 -0400
Fix artifact staging windows filename (#39916)
* Fix artifact staging filenames on Windows
* Fix tests in windows OS and exercises on windows test workflow
* Address comments
* Handle invalid chars in base as well
* Fix white spaces
---------
Co-authored-by: Atharv Urunkar <[email protected]>
---
.github/workflows/java_tests.yml | 29 +++--------
.../org/apache/beam/gradle/BeamModulePlugin.groovy | 5 ++
runners/java-fn-execution/build.gradle | 1 +
.../artifact/ArtifactStagingService.java | 19 ++++++--
.../fnexecution/environment/ProcessManager.java | 3 +-
.../artifact/ArtifactStagingServiceTest.java | 41 ++++++++++++++++
.../environment/ProcessManagerTest.java | 57 ++++++++++++++++------
runners/spark/job-server/spark_job_server.gradle | 2 +
8 files changed, 114 insertions(+), 43 deletions(-)
diff --git a/.github/workflows/java_tests.yml b/.github/workflows/java_tests.yml
index dcd6b052fde..1259904b7e2 100644
--- a/.github/workflows/java_tests.yml
+++ b/.github/workflows/java_tests.yml
@@ -63,38 +63,23 @@ jobs:
with:
gradle-command: test
arguments: -p sdks/java/core/
- - name: Upload test logs for :sdks:java:core:test
- uses: actions/upload-artifact@v7
- if: always()
- with:
- name: java_unit_tests-sdks-java-core-test-${{ matrix.os }}
- path: sdks/java/core/build/reports/tests/test
# :sdks:java:harness:test
- name: Run :sdks:java:harness:test
uses: ./.github/actions/gradle-command-self-hosted-action
with:
gradle-command: test
arguments: -p sdks/java/harness/
- if: always()
- - name: Upload test logs for :sdks:java:harness:test
- uses: actions/upload-artifact@v7
- if: always()
- with:
- name: java_unit_tests-sdks-java-harness-test-${{ matrix.os }}
- path: sdks/java/harness/build/reports/tests/test
- # :runners:core-java:test
- - name: Run :runners:core-java:test
+ # :runners:core-java:test and :runners:java-fn-execution:test
+ - name: Run basic runner tests
uses: ./.github/actions/gradle-command-self-hosted-action
with:
- gradle-command: test
- arguments: -p runners/core-java/
- if: always()
- - name: Upload test logs for :runners:core-java:test
+ gradle-command: :runners:core-java:test
:runners:java-fn-execution:test
+ - name: Upload test logs
uses: actions/upload-artifact@v7
- if: always()
+ if: ${{ !success() }}
with:
- name: java_unit_tests-runners-core-java-test-${{ matrix.os }}
- path: runners/core-java/build/reports/tests/test
+ name: java_unit_tests-${{ matrix.os }}
+ path: "**/build/reports/tests/"
java_wordcount_direct_runner:
name: 'Java Wordcount Direct Runner'
diff --git
a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
index be10dd4779c..d76180a5eec 100644
--- a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
+++ b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
@@ -1258,6 +1258,11 @@ class BeamModulePlugin implements Plugin<Project> {
useJUnit {}
// default maxHeapSize on gradle 5 is 512m, lets increase to handle
more demanding tests
maxHeapSize = '2g'
+ // Windows OS: Snappy needs an executable temp dir for native lib.
Default AppData/Temp
+ // failing with Access error without elevated permissions
+ if (System.getProperty("os.name").toLowerCase().contains("windows")) {
+ systemProperty 'org.xerial.snappy.tempdir',
System.getProperty('org.xerial.snappy.tempdir') ?:
"${project.rootDir.absolutePath}/build/snappy_bin"
+ }
}
// NOTE: Use the character class "[.]" instead of an escaped "\\." to
match a literal dot in
diff --git a/runners/java-fn-execution/build.gradle
b/runners/java-fn-execution/build.gradle
index ab45e10b208..2dd95d6d205 100644
--- a/runners/java-fn-execution/build.gradle
+++ b/runners/java-fn-execution/build.gradle
@@ -24,6 +24,7 @@ description = "Apache Beam :: Runners :: Java Fn Execution"
dependencies {
implementation library.java.vendored_guava_32_1_2_jre
+ implementation library.java.commons_lang3
implementation project(":runners:core-java")
compileOnly project(":sdks:java:harness")
implementation project(path: ":model:pipeline", configuration: "shadow")
diff --git
a/runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/artifact/ArtifactStagingService.java
b/runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/artifact/ArtifactStagingService.java
index 21c653ab6f1..512d435db32 100644
---
a/runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/artifact/ArtifactStagingService.java
+++
b/runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/artifact/ArtifactStagingService.java
@@ -17,6 +17,8 @@
*/
package org.apache.beam.runners.fnexecution.artifact;
+import static org.apache.commons.lang3.SystemUtils.IS_OS_WINDOWS;
+
import com.google.auto.value.AutoValue;
import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.io.IOException;
@@ -40,6 +42,7 @@ import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
+import java.util.regex.Pattern;
import org.apache.beam.model.jobmanagement.v1.ArtifactApi;
import org.apache.beam.model.jobmanagement.v1.ArtifactStagingServiceGrpc;
import org.apache.beam.model.pipeline.v1.RunnerApi;
@@ -58,7 +61,6 @@ import org.apache.beam.vendor.grpc.v1p69p0.io.grpc.Status;
import org.apache.beam.vendor.grpc.v1p69p0.io.grpc.StatusException;
import org.apache.beam.vendor.grpc.v1p69p0.io.grpc.stub.StreamObserver;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions;
-import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Splitter;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.hash.Hashing;
import org.slf4j.Logger;
@@ -72,6 +74,9 @@ public class ArtifactStagingService
private static final Logger LOG =
LoggerFactory.getLogger(ArtifactStagingService.class);
+ private static final Pattern WINDOWS_INVALID_CHARS =
+ Pattern.compile("[<>:\"/\\\\|?*\\x00-\\x1F]");
+
private final ArtifactDestinationProvider destinationProvider;
private final ConcurrentMap<String, Map<String,
List<RunnerApi.ArtifactInformation>>> toStage =
@@ -525,10 +530,16 @@ public class ArtifactStagingService
} catch (InvalidProtocolBufferException exn) {
throw new RuntimeException(exn);
}
- // Limit to the last contiguous alpha-numeric sequence. In particular,
this will exclude
+ // Limit to the last contiguous valid windows path chars. In
particular, this will exclude
// all path separators.
- List<String> components =
Splitter.onPattern("[^A-Za-z-_.]]").splitToList(path);
- String base = components.get(components.size() - 1);
+ String base =
+ WINDOWS_INVALID_CHARS
+ .splitAsStream(path)
+ .reduce((first, second) -> second)
+ .orElse("artifact");
+ if (IS_OS_WINDOWS) {
+ environment =
WINDOWS_INVALID_CHARS.matcher(environment).replaceAll("_");
+ }
return clip(
String.format("%s-%s-%s", idGenerator.getId(), clip(environment,
25), base), 100);
}
diff --git
a/runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/environment/ProcessManager.java
b/runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/environment/ProcessManager.java
index 86299762783..8fc407a0514 100644
---
a/runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/environment/ProcessManager.java
+++
b/runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/environment/ProcessManager.java
@@ -18,6 +18,7 @@
package org.apache.beam.runners.fnexecution.environment;
import static
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkNotNull;
+import static org.apache.commons.lang3.SystemUtils.IS_OS_WINDOWS;
import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.io.File;
@@ -113,7 +114,7 @@ public class ProcessManager {
} else {
// Pipe stdout and stderr to /dev/null to avoid blocking the process due
to filled PIPE
// buffer
- if (System.getProperty("os.name", "").startsWith("Windows")) {
+ if (IS_OS_WINDOWS) {
outputFile = new File("nul");
} else {
outputFile = new File("/dev/null");
diff --git
a/runners/java-fn-execution/src/test/java/org/apache/beam/runners/fnexecution/artifact/ArtifactStagingServiceTest.java
b/runners/java-fn-execution/src/test/java/org/apache/beam/runners/fnexecution/artifact/ArtifactStagingServiceTest.java
index abdce59458d..68b7c3da76f 100644
---
a/runners/java-fn-execution/src/test/java/org/apache/beam/runners/fnexecution/artifact/ArtifactStagingServiceTest.java
+++
b/runners/java-fn-execution/src/test/java/org/apache/beam/runners/fnexecution/artifact/ArtifactStagingServiceTest.java
@@ -185,6 +185,47 @@ public class ArtifactStagingServiceTest {
checkArtifacts(contentsList, staged.get("env2"));
}
+ @Test
+ public void testStageArtifactsWithInvalidFilenameCharacters()
+ throws InterruptedException, ExecutionException {
+ String environment = "0:ref_Environment_default";
+ List<String> contentsList = ImmutableList.of("artifact-content");
+
+ stagingService.registerJob(
+ "stagingToken",
+ ImmutableMap.of(
+ environment,
+ Lists.transform(contentsList,
FakeArtifactRetrievalService::resolvedArtifact)));
+
+ ArtifactStagingService.offer(new FakeArtifactRetrievalService(),
stagingStub, "stagingToken");
+
+ Map<String, List<RunnerApi.ArtifactInformation>> staged =
+ stagingService.getStagedArtifacts("stagingToken");
+
+ assertEquals(1, staged.size());
+ checkArtifacts(contentsList, staged.get(environment));
+ }
+
+ @Test
+ public void testStageFileArtifactWithAbsolutePath() throws Exception {
+ java.io.File source = tempFolder.newFile("real-artifact.bin");
+ java.nio.file.Files.write(
+ source.toPath(),
"payload".getBytes(java.nio.charset.StandardCharsets.UTF_8));
+ RunnerApi.ArtifactInformation fileArtifact =
+ RunnerApi.ArtifactInformation.newBuilder()
+ .setTypeUrn(ArtifactRetrievalService.FILE_ARTIFACT_URN)
+ .setTypePayload(
+ RunnerApi.ArtifactFilePayload.newBuilder()
+ .setPath(source.getAbsolutePath())
+ .build()
+ .toByteString())
+ .setRoleUrn("beam:artifact:role:pip_requirements_file:v1")
+ .build();
+ stagingService.registerJob("fileToken", ImmutableMap.of("env",
ImmutableList.of(fileArtifact)));
+ ArtifactStagingService.offer(retrievalService, stagingStub, "fileToken");
+ assertEquals(1, stagingService.getStagedArtifacts("fileToken").size());
+ }
+
@SuppressWarnings("InlineMeInliner") // inline `Strings.repeat()` - Java 11+
API only
@Test(timeout = 60_000)
public void testDestinationFailureFailsOfferInsteadOfHanging() throws
Exception {
diff --git
a/runners/java-fn-execution/src/test/java/org/apache/beam/runners/fnexecution/environment/ProcessManagerTest.java
b/runners/java-fn-execution/src/test/java/org/apache/beam/runners/fnexecution/environment/ProcessManagerTest.java
index 4074bf94943..5cad166858c 100644
---
a/runners/java-fn-execution/src/test/java/org/apache/beam/runners/fnexecution/environment/ProcessManagerTest.java
+++
b/runners/java-fn-execution/src/test/java/org/apache/beam/runners/fnexecution/environment/ProcessManagerTest.java
@@ -17,9 +17,10 @@
*/
package org.apache.beam.runners.fnexecution.environment;
+import static org.apache.commons.lang3.SystemUtils.IS_OS_WINDOWS;
import static org.hamcrest.MatcherAssert.assertThat;
import static org.hamcrest.Matchers.containsString;
-import static org.hamcrest.Matchers.is;
+import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertNull;
@@ -45,14 +46,33 @@ import org.junit.runners.JUnit4;
public class ProcessManagerTest {
@Rule public transient Timeout globalTimeout = Timeout.seconds(600);
+ private static String bash(int pos) {
+ if (IS_OS_WINDOWS) {
+ switch (pos) {
+ case 0:
+ return "cmd";
+ case 1:
+ return "/c";
+ }
+ } else {
+ switch (pos) {
+ case 0:
+ return "bash";
+ case 1:
+ return "-c";
+ }
+ }
+ throw new IllegalArgumentException(String.format("Unknown pos %d", pos));
+ }
+
@Test
public void testRunSimpleCommand() throws IOException {
ProcessManager processManager = ProcessManager.create();
- processManager.startProcess("1", "bash", Collections.emptyList());
+ processManager.startProcess("1", bash(0), Collections.emptyList());
processManager.stopProcess("1");
- processManager.startProcess("2", "bash", Arrays.asList("-c", "ls"));
+ processManager.startProcess("2", bash(0), Arrays.asList(bash(1), "ls"));
processManager.stopProcess("2");
- processManager.startProcess("1", "bash", Arrays.asList("-c", "ls", "-l",
"-a"));
+ processManager.startProcess("1", bash(0), Arrays.asList(bash(1), "ls",
"-l", "-a"));
processManager.stopProcess("1");
}
@@ -70,9 +90,9 @@ public class ProcessManagerTest {
@Test
public void testDuplicateId() throws IOException {
ProcessManager processManager = ProcessManager.create();
- processManager.startProcess("1", "bash", Arrays.asList("-c", "ls"));
+ processManager.startProcess("1", bash(0), Arrays.asList(bash(1), "ls"));
try {
- processManager.startProcess("1", "bash", Arrays.asList("-c", "ls"));
+ processManager.startProcess("1", bash(0), Arrays.asList(bash(1), "ls"));
fail();
} catch (IllegalStateException e) {
// this is what we want
@@ -85,7 +105,7 @@ public class ProcessManagerTest {
public void testLivenessCheck() throws IOException {
ProcessManager processManager = ProcessManager.create();
ProcessManager.RunningProcess process =
- processManager.startProcess("1", "bash", Arrays.asList("-c", "sleep",
"1000"));
+ processManager.startProcess("1", bash(0), Arrays.asList(bash(1),
"sleep", "1000"));
process.isAliveOrThrow();
processManager.stopProcess("1");
try {
@@ -102,14 +122,19 @@ public class ProcessManagerTest {
ProcessManager.RunningProcess process =
processManager.startProcess(
"1",
- "bash",
- Arrays.asList("-c", "sleep $PARAM"),
- Collections.singletonMap("PARAM", "-h"));
+ bash(0),
+ Arrays.asList(
+ bash(1), "exit " + (IS_OS_WINDOWS ? "%TEST_ENV_PARAM%" :
"$TEST_ENV_PARAM")),
+ Collections.singletonMap("TEST_ENV_PARAM", "42"));
for (int i = 0; i < 10 && process.getUnderlyingProcess().isAlive(); i++) {
Thread.sleep(100);
}
- assertThat(process.getUnderlyingProcess().exitValue(), is(1));
- processManager.stopProcess("1");
+ int exCode = process.getUnderlyingProcess().exitValue();
+ try {
+ assertEquals(42, exCode);
+ } finally {
+ processManager.stopProcess("1");
+ }
}
@Test
@@ -120,8 +145,8 @@ public class ProcessManagerTest {
ProcessManager.RunningProcess process =
processManager.startProcess(
"1",
- "bash",
- Arrays.asList("-c", "echo 'testing123'"),
+ bash(0),
+ Arrays.asList(bash(1), "echo 'testing123'"),
Collections.emptyMap(),
outputFile);
for (int i = 0; i < 10 && process.getUnderlyingProcess().isAlive(); i++) {
@@ -171,7 +196,7 @@ public class ProcessManagerTest {
assertNull(ProcessManager.shutdownHook);
processManager.startProcess(
- "1", "bash", Arrays.asList("-c", "echo 'testing123'"),
Collections.emptyMap());
+ "1", bash(0), Arrays.asList(bash(1), "echo 'testing123'"),
Collections.emptyMap());
// the shutdown hook will be created when process is started
assertNotNull(ProcessManager.shutdownHook);
// check the shutdown hook is registered
@@ -180,7 +205,7 @@ public class ProcessManagerTest {
Runtime.getRuntime().addShutdownHook(ProcessManager.shutdownHook);
processManager.startProcess(
- "2", "bash", Arrays.asList("-c", "echo 'testing123'"),
Collections.emptyMap());
+ "2", bash(0), Arrays.asList(bash(1), "echo 'testing123'"),
Collections.emptyMap());
processManager.stopProcess("1");
// the shutdown hook will be not removed if there are still processes alive
diff --git a/runners/spark/job-server/spark_job_server.gradle
b/runners/spark/job-server/spark_job_server.gradle
index 2811f875f84..0166f065401 100644
--- a/runners/spark/job-server/spark_job_server.gradle
+++ b/runners/spark/job-server/spark_job_server.gradle
@@ -232,6 +232,8 @@ def portableValidatesRunnerTask(String name, boolean
streaming, boolean docker,
"beam.spark.test.reuseSparkContext": "false",
"spark.ui.enabled": "false",
"spark.ui.showConsoleProgress": "false",
+ // For Windows OS: additional config needed as below. See
https://cwiki.apache.org/confluence/spaces/HADOOP2/pages/120730292/WindowsProblems
+ // 'hadoop.home.dir':
"${rootDir.absolutePath}\\build\\hadoop-home",
],
testCategories: testCategories,
testFilter: testFilter,