Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 5 additions & 4 deletions .github/workflows/beam_PostRelease_NightlySnapshot.yml
Original file line number Diff line number Diff line change
Expand Up @@ -19,12 +19,13 @@ on:
workflow_dispatch:
inputs:
RELEASE:
description: Beam version of current release (e.g. 2.XX.0)
required: true
description: Beam version of current release (pass in empty string for nightly SNAPSHOT)
required: false
default: '2.XX.0'
SNAPSHOT_URL:
description: Location of the staged artifacts in Maven central (https://repository.apache.org/content/repositories/orgapachebeam-NNNN/).
required: true
description: Location of the staged artifacts in Maven central (https://repository.apache.org/content/repositories/orgapachebeam-NNNN/ or leave empty for snapshots).
required: false
default: ''
schedule:
- cron: '15 16 * * *'

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2643,7 +2643,7 @@ class BeamModulePlugin implements Plugin<Project> {
JavaExamplesArchetypeValidationConfiguration config = it as JavaExamplesArchetypeValidationConfiguration

def taskName = "run${config.type}Java${config.runner}"
def releaseVersion = project.findProperty('ver') ?: project.version
def releaseVersion = (project.findProperty('ver') && project.findProperty('ver') != '2.XX.0') ? project.findProperty('ver') : project.version

@Abacn Abacn Aug 7, 2026

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: Logically not recommending changing the common BeamModulePlugin.groovy to accommodate a specific workflow, beam_PostRelease_NightlySnapshot.yml‎. Instead, we can change the workflow like

-Pver='${{ (github.event.inputs.RELEASE == "2.xx.0" || '') || github.event.inputs.RELEASE }}'

def releaseRepo = project.findProperty('repourl') ?: 'https://repository.apache.org/content/repositories/snapshots'
// shared maven local path for maven archetype projects
def sharedMavenLocal = project.findProperty('mavenLocalPath') ?: ''
Expand Down
2 changes: 1 addition & 1 deletion release/build.gradle.kts
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ task("runJavaExamplesValidationTask") {
dependsOn(":runners:spark:3:runQuickstartJavaSpark")
dependsOn(":runners:flink:2.2:runQuickstartJavaFlinkLocal")
dependsOn(":runners:direct-java:runMobileGamingJavaDirect")
if (project.hasProperty("ver") || !project.version.toString().endsWith("SNAPSHOT")) {
if ((project.findProperty("ver")?.toString()?.isNotEmpty() == true && project.findProperty("ver") != "2.XX.0") || !project.version.toString().endsWith("SNAPSHOT")) {
// only run one variant of MobileGaming on Dataflow for nightly
dependsOn(":runners:google-cloud-dataflow-java:runMobileGamingJavaDataflow")
}
Expand Down
156 changes: 135 additions & 21 deletions release/src/main/groovy/TestScripts.groovy
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,10 @@ import groovy.util.CliBuilder
*/
class TestScripts {

class BackgroundProcessInfo {
String cmd
}

// Global state to maintain when running the steps
class var {
static File startDir
Expand All @@ -37,6 +41,8 @@ class TestScripts {
static String bqDataset
static String pubsubTopic
static String mavenLocalPath
static List<Process> backgroundProcesses = Collections.synchronizedList(new ArrayList<Process>())
static Map<Process, BackgroundProcessInfo> backgroundProcessInfo = Collections.synchronizedMap(new HashMap<Process, BackgroundProcessInfo>())
}

def TestScripts(String[] args) {
Expand Down Expand Up @@ -79,6 +85,10 @@ class TestScripts {
var.mavenLocalPath = options.mavenLocalPath
println "Maven local path: ${var.mavenLocalPath}"
}

Runtime.getRuntime().addShutdownHook(new Thread({
stopAllBackgroundProcesses()
}))
}

def ver() {
Expand Down Expand Up @@ -135,6 +145,75 @@ class TestScripts {
}
}

// Run a command in the background, returning the Process object.
public Process runBackground(String cmd) {
println cmd
if (cmd.startsWith("mvn ")) {
return _mvnBackground(cmd.substring(4))
} else {
return _executeBackground(cmd)
}
}

// Check whether any background processes exited unexpectedly with a non-zero exit code
public void checkBackgroundProcesses() {
def procs = new ArrayList<>(var.backgroundProcesses)
for (Process proc : procs) {
if (proc != null && !proc.isAlive()) {
int exitVal = proc.exitValue()
if (exitVal != 0) {
def info = var.backgroundProcessInfo.get(proc)
String cmd = info ? info.cmd : "unknown command"
error("Background command failed with exit code ${exitVal}: ${cmd}")
}
}
}
}

// Stop/kill a background process and all its descendants.
public void stopProcess(Process proc) {
if (proc != null) {
if (!proc.isAlive()) {
int exitVal = proc.exitValue()
var.backgroundProcesses.remove(proc)
def info = var.backgroundProcessInfo.remove(proc)
if (exitVal != 0) {
String cmd = info ? info.cmd : "unknown command"
error("Background command failed with exit code ${exitVal}: ${cmd}")
}
} else {
try {
proc.descendants().forEach { it.destroyForcibly() }
} catch (Throwable ignored) {
}
proc.destroyForcibly()
proc.waitFor(10, java.util.concurrent.TimeUnit.SECONDS)
var.backgroundProcesses.remove(proc)
var.backgroundProcessInfo.remove(proc)
}
}
}

// Stop all active background processes.
public void stopAllBackgroundProcesses() {
def procs = new ArrayList<>(var.backgroundProcesses)
procs.each { proc ->
if (proc != null && proc.isAlive()) {
try {
proc.descendants().forEach { it.destroyForcibly() }
} catch (Throwable ignored) {
}
proc.destroyForcibly()
try {
proc.waitFor(10, java.util.concurrent.TimeUnit.SECONDS)
} catch (Throwable ignored) {
}
}
var.backgroundProcesses.remove(proc)
var.backgroundProcessInfo.remove(proc)
}
}

// Check for expected results in actual stdout from previous command, if fails, log errors then exit.
public void see(String expected, String actual) {
if (!actual.contains(expected)) {
Expand All @@ -159,13 +238,16 @@ class TestScripts {

// Cleanup and print success
public void done() {
checkBackgroundProcesses()
stopAllBackgroundProcesses()
var.startDir.deleteDir()
println "[SUCCESS]"
System.exit(0)
}

// Run a single command, capture output, verify return code is 0
private String _execute(String cmd) {
checkBackgroundProcesses()
def shell = "sh -c cmd".split(' ')
shell[2] = cmd
def pb = new ProcessBuilder(shell)
Expand All @@ -187,6 +269,27 @@ class TestScripts {
return output_text
}

// Run a single command asynchronously in the background
private Process _executeBackground(String cmd) {
def shell = "sh -c cmd".split(' ')
shell[2] = cmd
def pb = new ProcessBuilder(shell)
pb.directory(var.curDir)
pb.redirectErrorStream(true)
def proc = pb.start()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could we retain and check background process exit status here? _executeBackground only drains stdout, and stopProcess removes an already-exited process without inspecting its result. DirectRunner keeps leaderboard_DirectRunner_* tables, so failed Injector or LeaderBoard commands can match stale rows and report [SUCCESS]. Reproduced with both background Maven commands exiting 42 while script exited 0.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks - addressed this + did some other cleanup (running https://github.com/apache/beam/actions/runs/31172410588 to verify I didn't break anything with this)

var.backgroundProcesses.add(proc)
var.backgroundProcessInfo.put(proc, new BackgroundProcessInfo(cmd: cmd))
Thread.startDaemon {
try {
proc.inputStream.eachLine {
println it
}
} catch (Throwable ignored) {
}
}
return proc
}

// Change directory
private void _chdir(String subdir) {
var.curDir = new File(var.curDir.absolutePath, subdir)
Expand All @@ -195,46 +298,57 @@ class TestScripts {
}
}

// Run a maven command, setting up a new local repository and a settings.xml with a custom repository if needed
private String _mvn(String args) {
// Build the maven command string with custom repository and settings.xml
private String _buildMvnCmd(String args) {
String mvnlocalPath = var.mavenLocalPath
if (!(var.mavenLocalPath)) {
mvnlocalPath = var.startDir
}
def m2 = new File(mvnlocalPath, ".m2/repository")
m2.mkdirs()
def settings = new File(mvnlocalPath, "settings.xml")
if(!settings.exists()) {
settings.write """
<settings>
<localRepository>${m2.absolutePath}</localRepository>
<profiles>
<profile>
<id>testrel</id>
<repositories>
<repository>
<id>test.release</id>
<url>${var.repoUrl}</url>
</repository>
</repositories>
</profile>
</profiles>
</settings>
"""
if (!settings.exists()) {
settings.write """
<settings>
<localRepository>${m2.absolutePath}</localRepository>
<profiles>
<profile>
<id>testrel</id>
<repositories>
<repository>
<id>test.release</id>
<url>${var.repoUrl}</url>
</repository>
</repositories>
</profile>
</profiles>
</settings>
"""
}
def cmd = "mvn ${args} -s ${settings.absolutePath} -Ptestrel -B"
String path = System.getenv("PATH");
String path = System.getenv("PATH")
// Set the path on jenkins executors to use a recent maven
// MAVEN_HOME is not set on some executors, so default to 3.5.2
String maven_home = System.getenv("MAVEN_HOME") ?: '/usr/local/maven'
println "Using maven ${maven_home}"
def mvnPath = "${maven_home}/bin"
def setPath = "export PATH=\"${mvnPath}:${path}\" && "
return _execute(setPath + cmd)
return setPath + cmd
}

// Run a maven command, setting up a new local repository and a settings.xml with a custom repository if needed
private String _mvn(String args) {
return _execute(_buildMvnCmd(args))
}

// Run a maven command in the background
private Process _mvnBackground(String args) {
return _executeBackground(_buildMvnCmd(args))
}

// Clean up and report error
public void error(String text) {
stopAllBackgroundProcesses()
var.startDir.deleteDir()
println "[ERROR] $text"
System.exit(1)
Expand Down
75 changes: 49 additions & 26 deletions release/src/main/groovy/mobilegaming-java-dataflow.groovy
Original file line number Diff line number Diff line change
Expand Up @@ -120,37 +120,53 @@ class LeaderBoardRunner {
].join(",")

String tables = t.run("bq query --use_legacy_sql=false 'SELECT table_name FROM ${dataset}.INFORMATION_SCHEMA.TABLES'")

if (!tables.contains(userTable)) {
t.intent("Creating table: ${userTable}")
t.run("bq mk --table ${dataset}.${userTable} ${userSchema}")
if (tables.contains(userTable)) {
t.run("bq rm -f -t ${dataset}.${userTable}")
}
if (tables.contains(teamTable)) {
t.run("bq rm -f -t ${dataset}.${teamTable}")
}
if (!tables.contains(teamTable)) {
t.intent("Creating table: ${teamTable}")
t.run("bq mk --table ${dataset}.${teamTable} ${teamSchema}")
int retries = 10
boolean deleted = false
for (int i = 0; i < retries; i++) {
tables = t.run("bq query --use_legacy_sql=false 'SELECT table_name FROM ${dataset}.INFORMATION_SCHEMA.TABLES'")
if (!tables.contains(userTable) && !tables.contains(teamTable)) {
deleted = true
break
}
sleep(3000)
}
if (!deleted) {
t.error("Timed out waiting for tables ${userTable} / ${teamTable} to be deleted.")
}

t.intent("Creating table: ${userTable}")
t.run("bq mk --table ${dataset}.${userTable} ${userSchema}")
t.intent("Creating table: ${teamTable}")
t.run("bq mk --table ${dataset}.${teamTable} ${teamSchema}")

// Verify that the tables have been created successfully
tables = t.run("bq query --use_legacy_sql=false 'SELECT table_name FROM ${dataset}.INFORMATION_SCHEMA.TABLES'")
while (!tables.contains(userTable) || !tables.contains(teamTable)) {
sleep(3000)
boolean created = false
for (int i = 0; i < retries; i++) {
tables = t.run("bq query --use_legacy_sql=false 'SELECT table_name FROM ${dataset}.INFORMATION_SCHEMA.TABLES'")
if (tables.contains(userTable) && tables.contains(teamTable)) {
created = true
break
}
sleep(3000)
}
if (!created) {
t.error("Timed out waiting for tables ${userTable} / ${teamTable} to be created.")
}
println "Tables ${userTable} and ${teamTable} created successfully."

def InjectorThread = Thread.start() {
t.run(mobileGamingCommands.createInjectorCommand())
}
def injectorProcess = t.runBackground(mobileGamingCommands.createInjectorCommand())

String jobName = "leaderboard-validation-" + new Date().getTime() + "-" + new Random().nextInt(1000)
def LeaderBoardThread = Thread.start() {
if (useStreamingEngine) {
t.run(mobileGamingCommands.createPipelineCommand(
"LeaderBoardWithStreamingEngine", runner, jobName, "LeaderBoard"))
} else {
t.run(mobileGamingCommands.createPipelineCommand("LeaderBoard", runner, jobName))
}
}
def leaderBoardProcess = useStreamingEngine ?
t.runBackground(mobileGamingCommands.createPipelineCommand(
"LeaderBoardWithStreamingEngine", runner, jobName, "LeaderBoard")) :
t.runBackground(mobileGamingCommands.createPipelineCommand("LeaderBoard", runner, jobName))

t.run("gcloud dataflow jobs list | grep pyflow-wordstream-candidate | grep Running | cut -d' ' -f1")

Expand All @@ -175,8 +191,8 @@ class LeaderBoardRunner {
println "Waiting for pipeline to produce more results..."
sleep(60000) // wait for 1 min
}
InjectorThread.stop()
LeaderBoardThread.stop()
t.stopProcess(injectorProcess)
t.stopProcess(leaderBoardProcess)
t.run("""RUNNING_JOB=`gcloud dataflow jobs list | grep ${jobName} | grep Running | cut -d' ' -f1`
if [ ! -z "\${RUNNING_JOB}" ]
then
Expand All @@ -202,10 +218,17 @@ fi

// It will take couple seconds to clean up tables.
// This loop makes sure tables are completely deleted before running the pipeline
tables = t.run("bq query --use_legacy_sql=false 'SELECT table_name FROM ${dataset}.INFORMATION_SCHEMA.TABLES'")
while (tables.contains(userTable) || tables.contains(teamTable)) {
sleep(3000)
deleted = false
for (int i = 0; i < retries; i++) {
tables = t.run("bq query --use_legacy_sql=false 'SELECT table_name FROM ${dataset}.INFORMATION_SCHEMA.TABLES'")
if (!tables.contains(userTable) && !tables.contains(teamTable)) {
deleted = true
break
}
sleep(3000)
}
if (!deleted) {
println "Warning: Timed out waiting for tables ${userTable} / ${teamTable} to be deleted."
}
}
}
Expand Down
Loading
Loading