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

yuqi1129 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gravitino.git


The following commit(s) were added to refs/heads/main by this push:
     new 1c70cdcc6a [#13372] fix(core): Make job staging directories 
independent of template and metalake names (#13425)
1c70cdcc6a is described below

commit 1c70cdcc6acfea1b816684864b1784d34421b80a
Author: Jerry Shao <[email protected]>
AuthorDate: Wed Sep 23 21:52:09 2026 +0800

    [#13372] fix(core): Make job staging directories independent of template 
and metalake names (#13425)
    
    ### What changes were proposed in this pull request?
    
    A job's staging directory is now
    `<gravitino.job.stagingDir>/job-runs/job-<id>`, derived from the job id
    alone, instead of `<stagingDir>/<metalake>/<template>/job-<id>`.
    
    Both `cleanUpStagingDirs()` and `deleteJobTemplate()` delete a job's
    directory through one helper. It first tries the new location. If that
    doesn't exist, the job was submitted by an earlier version, and its
    directory is looked up by id under `<stagingDir>/*/*/job-<id>` instead
    of being rebuilt from the current names. This also removes directories
    that were already leaked by an earlier rename. Symbolic links and hidden
    entries such as `.job-output-index` are skipped, and an unreadable
    directory (e.g. `lost+found`) doesn't stop the lookup.
    
    Also updates the staging paths in the table maintenance docs, and fixes
    the `JobIT` check for a rejected job's staging directory, which passed
    vacuously with the new layout.
    
    ### Why are the changes needed?
    
    The cleanup rebuilt the staging path from the current template and
    metalake names. After a template or metalake rename, the path didn't
    exist, nothing was deleted, and the job entity was removed anyway, so
    the directory (executables, scripts, jars, output logs) leaked forever.
    The metalake rename case is verified by a new test.
    
    Fix: #13372
    
    ### Does this PR introduce _any_ user-facing change?
    
    The layout of the job staging directory changes as described above.
    Directories of jobs submitted before the upgrade are still cleaned up.
    No API or configuration changes.
    
    During a rolling upgrade with a shared staging directory, a server
    running the old version can't find the directories of jobs submitted by
    upgraded servers, so the jobs it cleans up leave their directories
    behind.
    
    ### How was this patch tested?
    
    - New tests in `TestJobManagerMultiNode` (real JDBC backend and
    `LocalJobExecutor`) for template rename + expiry, repeated template
    renames + `deleteJobTemplate`, and metalake rename + expiry. All three
    fail on main.
    - New unit tests in `TestJobManager` for the legacy directory lookup,
    including symlinks, hidden directories, same-name files and unreadable
    directories.
    - `./gradlew :core:test -PskipITs --tests 'org.apache.gravitino.job.*'`,
    and `JobIT` in embedded mode.
    
    🤖 Generated with [Claude Code](https://claude.com/claude-code)
    
    ---------
    
    Co-authored-by: Claude Opus 5 <[email protected]>
---
 .../gravitino/client/integration/test/JobIT.java   |  16 +-
 .../java/org/apache/gravitino/job/JobManager.java  | 140 +++++++++---
 .../org/apache/gravitino/job/TestJobManager.java   | 237 +++++++++++++++++++--
 .../gravitino/job/TestJobManagerMultiNode.java     |  90 +++++++-
 .../gravitino/job/local/TestLocalJobExecutor.java  |  13 +-
 .../optimizer-cli-reference.md                     |   2 +-
 .../optimizer-troubleshooting.md                   |   2 +-
 7 files changed, 441 insertions(+), 59 deletions(-)

diff --git 
a/clients/client-java/src/test/java/org/apache/gravitino/client/integration/test/JobIT.java
 
b/clients/client-java/src/test/java/org/apache/gravitino/client/integration/test/JobIT.java
index 08f38a9d24..1442354417 100644
--- 
a/clients/client-java/src/test/java/org/apache/gravitino/client/integration/test/JobIT.java
+++ 
b/clients/client-java/src/test/java/org/apache/gravitino/client/integration/test/JobIT.java
@@ -19,6 +19,7 @@
 package org.apache.gravitino.client.integration.test;
 
 import com.google.common.collect.ImmutableMap;
+import com.google.common.collect.ImmutableSet;
 import com.google.common.collect.Lists;
 import java.io.File;
 import java.io.IOException;
@@ -380,6 +381,8 @@ public class JobIT extends BaseIT {
             .withClassName("org.apache.gravitino.test.SparkJob")
             .build();
     Assertions.assertDoesNotThrow(() -> 
metalake.registerJobTemplate(template));
+    File jobRunsDir = new File(testStagingDir, "job-runs");
+    Set<String> jobStagingDirsBefore = listFileNames(jobRunsDir);
 
     // The run request is rejected with the reason instead of being queued and 
failing later.
     IllegalArgumentException e =
@@ -396,9 +399,11 @@ public class JobIT extends BaseIT {
 
     // No job is created, and the staging directory of the rejected job is 
removed.
     Assertions.assertTrue(metalake.listJobs(template.name()).isEmpty());
-    String[] jobStagingDirs =
-        new File(testStagingDir, METALAKE_NAME + File.separator + 
template.name()).list();
-    Assertions.assertTrue(jobStagingDirs == null || jobStagingDirs.length == 
0);
+    // Directories of earlier jobs may expire meanwhile, but no new one may be 
left behind.
+    Set<String> jobStagingDirsAfter = listFileNames(jobRunsDir);
+    Assertions.assertTrue(
+        jobStagingDirsBefore.containsAll(jobStagingDirsAfter),
+        "Staging directories left behind: " + jobStagingDirsAfter);
 
     // Shell jobs are not affected by the missing Spark installation.
     JobTemplate shellTemplate = 
builder.withName("test_run_shell_without_spark_submit").build();
@@ -671,6 +676,11 @@ public class JobIT extends BaseIT {
     return job;
   }
 
+  private static Set<String> listFileNames(File dir) {
+    String[] names = dir.list();
+    return names == null ? Collections.emptySet() : ImmutableSet.copyOf(names);
+  }
+
   private String generateTestEntryScript() {
     String content =
         "#!/bin/bash\n"
diff --git a/core/src/main/java/org/apache/gravitino/job/JobManager.java 
b/core/src/main/java/org/apache/gravitino/job/JobManager.java
index 7e67602a43..a2f2069d08 100644
--- a/core/src/main/java/org/apache/gravitino/job/JobManager.java
+++ b/core/src/main/java/org/apache/gravitino/job/JobManager.java
@@ -27,7 +27,11 @@ import com.google.common.base.Preconditions;
 import java.io.File;
 import java.io.IOException;
 import java.net.URI;
+import java.nio.file.DirectoryIteratorException;
+import java.nio.file.DirectoryStream;
 import java.nio.file.Files;
+import java.nio.file.LinkOption;
+import java.nio.file.Path;
 import java.time.Instant;
 import java.util.ArrayList;
 import java.util.Arrays;
@@ -84,14 +88,20 @@ public class JobManager implements JobOperationDispatcher {
 
   private static final Pattern PLACEHOLDER_PATTERN = 
Pattern.compile("\\{\\{([\\w.-]+)\\}\\}");
 
-  private static final String JOB_STAGING_DIR =
-      File.separator
-          + "%s"
-          + File.separator
-          + "%s"
-          + File.separator
-          + JobHandle.JOB_ID_PREFIX
-          + "%s";
+  // A job's staging directory is <stagingDir>/job-runs/job-<id>, derived from 
the job id alone:
+  // template and metalake names can change after the job runs, so a path 
built from them can't be
+  // rebuilt at cleanup time. The name contains '-', which metalake names 
can't, so it never
+  // collides with a metalake directory of the legacy 
<stagingDir>/<metalake>/<template>/job-<id>
+  // layout.
+  private static final String JOB_RUNS_DIR_NAME = "job-runs";
+
+  private static final Pattern JOB_DIR_PATTERN =
+      Pattern.compile(Pattern.quote(JobHandle.JOB_ID_PREFIX) + "\\d+");
+
+  // Bounds how deep a legacy staging directory is looked up below a job 
template's first path
+  // element. A template name nested deeper than this keeps its directory, 
which is better than
+  // walking an unrelated directory tree an operator put in the staging 
directory.
+  private static final int LEGACY_JOB_STAGING_DIR_MAX_DEPTH = 16;
 
   private static final long JOB_STAGING_DIR_CLEANUP_MIN_TIME_IN_MS = 600 * 
1000L; // 10 minute
 
@@ -312,16 +322,13 @@ public class JobManager implements JobOperationDispatcher 
{
     }
 
     // Only remove directories belonging to the observed jobs. A same-name 
template can be
-    // recreated after the metadata transaction commits, so its parent 
directory is not ours to
-    // delete.
+    // recreated after the metadata transaction commits, so a legacy template 
directory is not ours
+    // to delete.
     for (JobEntity job : jobs) {
-      String jobStagingPath =
-          stagingDir.getAbsolutePath()
-              + String.format(JOB_STAGING_DIR, metalake, 
job.jobTemplateName(), job.id());
       try {
-        FileUtils.deleteDirectory(new File(jobStagingPath));
+        deleteJobStagingDir(job);
       } catch (IOException e) {
-        LOG.error("Failed to delete job staging directory: {}", 
jobStagingPath, e);
+        LOG.error("Failed to delete the staging directory of job {}", 
job.name(), e);
       }
     }
 
@@ -496,10 +503,7 @@ public class JobManager implements JobOperationDispatcher {
 
     // Create staging directory.
     long jobId = idGenerator.nextId();
-    String jobStagingPath =
-        stagingDir.getAbsolutePath()
-            + String.format(JOB_STAGING_DIR, metalake, jobTemplateName, jobId);
-    File jobStagingDir = new File(jobStagingPath);
+    File jobStagingDir = jobStagingDir(jobId);
     try {
       Files.createDirectories(jobStagingDir.toPath());
     } catch (IOException e) {
@@ -821,15 +825,7 @@ public class JobManager implements JobOperationDispatcher {
             try {
               entityStore.delete(
                   NameIdentifierUtil.ofJob(metalake, job.name()), 
Entity.EntityType.JOB);
-
-              String jobStagingPath =
-                  stagingDir.getAbsolutePath()
-                      + String.format(JOB_STAGING_DIR, metalake, 
job.jobTemplateName(), job.id());
-              File jobStagingDir = new File(jobStagingPath);
-              if (jobStagingDir.exists()) {
-                FileUtils.deleteDirectory(jobStagingDir);
-                LOG.info("Deleted job staging directory {} for job {}", 
jobStagingPath, job.name());
-              }
+              deleteJobStagingDir(job);
             } catch (OptimisticLockException e) {
               // Keep the files when deletion loses its CAS. The next cleanup 
run re-reads the
               // job and checks retention eligibility again; this batch can 
process other jobs.
@@ -1108,6 +1104,94 @@ public class JobManager implements 
JobOperationDispatcher {
         .build();
   }
 
+  @VisibleForTesting
+  File jobStagingDir(long jobId) {
+    return new File(new File(stagingDir, JOB_RUNS_DIR_NAME), 
JobHandle.JOB_ID_PREFIX + jobId);
+  }
+
+  /**
+   * Finds the staging directory of a job submitted by an earlier Gravitino 
version, which used the
+   * {@code <stagingDir>/<metalake>/<template>/job-<id>} layout. It's looked 
up by the job id rather
+   * than rebuilt from the current names, because the template or metalake may 
have been renamed
+   * since the job ran. The job id is unique, so there is at most one match in 
practice. Called only
+   * for a job without a directory in the current layout, which after an 
upgrade are the jobs of the
+   * earlier version until they expire.
+   */
+  @VisibleForTesting
+  List<File> findLegacyJobStagingDirs(long jobId) throws IOException {
+    String jobDirName = JobHandle.JOB_ID_PREFIX + jobId;
+    List<File> legacyJobStagingDirs = new ArrayList<>();
+    try (DirectoryStream<Path> metalakeDirs = 
Files.newDirectoryStream(stagingDir.toPath())) {
+      for (Path metalakeDir : metalakeDirs) {
+        String name = metalakeDir.getFileName().toString();
+        // Metalake names can't start with '.', so hidden entries, e.g. the 
job output index
+        // directory, are not metalake directories. Symbolic links are never 
followed, so nothing
+        // outside the staging directory can be deleted.
+        if (name.equals(JOB_RUNS_DIR_NAME)
+            || name.startsWith(".")
+            || !Files.isDirectory(metalakeDir, LinkOption.NOFOLLOW_LINKS)) {
+          continue;
+        }
+
+        // An unreadable directory, e.g. "lost+found" when the staging 
directory is the root of a
+        // file system, must not prevent finding the job under the other 
directories.
+        try (DirectoryStream<Path> templateDirs = 
Files.newDirectoryStream(metalakeDir)) {
+          for (Path templateDir : templateDirs) {
+            if (Files.isDirectory(templateDir, LinkOption.NOFOLLOW_LINKS)) {
+              collectLegacyJobStagingDirs(
+                  templateDir, jobDirName, LEGACY_JOB_STAGING_DIR_MAX_DEPTH, 
legacyJobStagingDirs);
+            }
+          }
+        } catch (IOException | DirectoryIteratorException e) {
+          LOG.warn(
+              "Failed to look up the legacy staging directory of job {} under 
{}, skip it",
+              jobDirName,
+              metalakeDir,
+              e);
+        }
+      }
+    }
+    return legacyJobStagingDirs;
+  }
+
+  // The legacy layout put the template name in the path as is, and a template 
name may contain
+  // '/', so the job directory can be nested deeper than 
<metalake>/<template>/job-<id>, e.g. under
+  // a template named "team/etl". The staging directory of another job is 
never descended into: its
+  // contents are that job's own files.
+  private static void collectLegacyJobStagingDirs(
+      Path dir, String jobDirName, int remainingDepth, List<File> 
legacyJobStagingDirs)
+      throws IOException {
+    try (DirectoryStream<Path> children = Files.newDirectoryStream(dir)) {
+      for (Path child : children) {
+        if (!Files.isDirectory(child, LinkOption.NOFOLLOW_LINKS)) {
+          continue;
+        }
+
+        String name = child.getFileName().toString();
+        if (name.equals(jobDirName)) {
+          legacyJobStagingDirs.add(child.toFile());
+        } else if (remainingDepth > 0 && 
!JOB_DIR_PATTERN.matcher(name).matches()) {
+          collectLegacyJobStagingDirs(child, jobDirName, remainingDepth - 1, 
legacyJobStagingDirs);
+        }
+      }
+    }
+  }
+
+  private void deleteJobStagingDir(JobEntity job) throws IOException {
+    File jobStagingDir = jobStagingDir(job.id());
+    if (jobStagingDir.exists()) {
+      FileUtils.deleteDirectory(jobStagingDir);
+      LOG.info("Deleted job staging directory {} for job {}", jobStagingDir, 
job.name());
+      return;
+    }
+
+    for (File legacyJobStagingDir : findLegacyJobStagingDirs(job.id())) {
+      FileUtils.deleteDirectory(legacyJobStagingDir);
+      LOG.info(
+          "Deleted legacy job staging directory {} for job {}", 
legacyJobStagingDir, job.name());
+    }
+  }
+
   private void deleteStagingDirOfUnsubmittedJob(File jobStagingDir, long 
jobId) {
     // The job is not tracked by any job entity, so the periodic cleanup will 
never remove its
     // staging directory. A cleanup failure must not mask the original 
submission failure.
diff --git a/core/src/test/java/org/apache/gravitino/job/TestJobManager.java 
b/core/src/test/java/org/apache/gravitino/job/TestJobManager.java
index 1091b38fec..d19c274c02 100644
--- a/core/src/test/java/org/apache/gravitino/job/TestJobManager.java
+++ b/core/src/test/java/org/apache/gravitino/job/TestJobManager.java
@@ -33,6 +33,7 @@ import static org.mockito.Mockito.when;
 
 import com.google.common.collect.ImmutableList;
 import com.google.common.collect.ImmutableMap;
+import com.google.common.collect.ImmutableSet;
 import com.google.common.collect.Lists;
 import com.sun.net.httpserver.HttpServer;
 import java.io.File;
@@ -53,6 +54,7 @@ import java.util.UUID;
 import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.TimeUnit;
 import java.util.function.Function;
+import java.util.stream.Collectors;
 import javax.annotation.Nullable;
 import org.apache.commons.io.FileUtils;
 import org.apache.commons.lang3.ArrayUtils;
@@ -405,10 +407,7 @@ public class TestJobManager {
         .delete(
             NameIdentifierUtil.ofJobTemplate(metalake, "shell_job"),
             Entity.EntityType.JOB_TEMPLATE);
-    File directory =
-        new File(
-            testStagingDir,
-            metalake + File.separator + "shell_job" + File.separator + 
finishedJob.name());
+    File directory = jobManager.jobStagingDir(finishedJob.id());
     Assertions.assertTrue(directory.mkdirs() || directory.isDirectory());
     File artifact = new File(directory, "artifact");
     Assertions.assertTrue(artifact.createNewFile());
@@ -428,9 +427,7 @@ public class TestJobManager {
   @Test
   public void testDeletePreservesReplacementStaging() throws IOException {
     doReturn(Collections.emptyList()).when(jobManager).listJobs(metalake, 
Optional.of("shell_job"));
-    File replacementDir =
-        new File(
-            testStagingDir, metalake + File.separator + "shell_job" + 
File.separator + "job_999");
+    File replacementDir = jobManager.jobStagingDir(999L);
     File replacementArtifact = new File(replacementDir, "new-job-artifact");
     when(entityStore.delete(
             NameIdentifierUtil.ofJobTemplate(metalake, "shell_job"),
@@ -455,10 +452,7 @@ public class TestJobManager {
     doReturn(Collections.singletonList(finishedJob))
         .when(jobManager)
         .listJobs(metalake, Optional.of("shell_job"));
-    File directory =
-        new File(
-            testStagingDir,
-            metalake + File.separator + "shell_job" + File.separator + 
finishedJob.name());
+    File directory = jobManager.jobStagingDir(finishedJob.id());
     Assertions.assertTrue(directory.mkdirs());
     File artifact = new File(directory, "artifact");
     Assertions.assertTrue(artifact.createNewFile());
@@ -824,10 +818,9 @@ public class TestJobManager {
 
     // No job entity is registered and the staging directory of the rejected 
job is removed.
     verify(entityStore, never()).put(any(JobEntity.class), anyBoolean());
-    File templateStagingDir =
-        new File(testStagingDir, metalake + File.separator + 
shellJobTemplate.name());
-    String[] jobStagingDirs = templateStagingDir.list();
-    Assertions.assertTrue(jobStagingDirs == null || jobStagingDirs.length == 
0);
+    File jobRunsDir = jobManager.jobStagingDir(0L).getParentFile();
+    Assertions.assertTrue(jobRunsDir.isDirectory(), "The job staging directory 
was never created");
+    Assertions.assertArrayEquals(new String[0], jobRunsDir.list());
   }
 
   @Test
@@ -1639,7 +1632,7 @@ public class TestJobManager {
     for (JobEntity job : ImmutableList.of(queuedJob, startedJob, 
cancellingJob, activeJob)) {
       stubEntityStoreUpdateToApply(job, job);
     }
-    File jobStagingDir = new File(testStagingDir, metalake + "/shell_job/" + 
startedJob.name());
+    File jobStagingDir = jobManager.jobStagingDir(startedJob.id());
     Assertions.assertTrue(jobStagingDir.mkdirs());
 
     long beforeCleanUp = System.currentTimeMillis();
@@ -1744,8 +1737,8 @@ public class TestJobManager {
         .thenThrow(new OptimisticLockException("job changed"))
         .thenReturn(true);
     when(entityStore.delete(otherIdent, 
Entity.EntityType.JOB)).thenReturn(true);
-    File conflictedDir = new File(testStagingDir, metalake + "/shell_job/" + 
conflicted.name());
-    File otherDir = new File(testStagingDir, metalake + "/shell_job/" + 
other.name());
+    File conflictedDir = jobManager.jobStagingDir(conflicted.id());
+    File otherDir = jobManager.jobStagingDir(other.id());
     Assertions.assertTrue(conflictedDir.mkdirs());
     Assertions.assertTrue(otherDir.mkdirs());
     File artifact = new File(conflictedDir, "artifact");
@@ -1760,6 +1753,214 @@ public class TestJobManager {
     verify(entityStore, times(1)).delete(otherIdent, Entity.EntityType.JOB);
   }
 
+  @Test
+  public void testCleanUpStagingDirsDeletesLegacyStagingDirOfRenamedTemplate() 
throws IOException {
+    // The job ran from template "shell_job" before an upgrade, in the legacy 
layout, and the
+    // template was renamed afterwards, so the job now reports the new name.
+    JobEntity job = expiredJob();
+    JobEntity renamedJob =
+        JobEntity.builder()
+            .withId(job.id())
+            .withJobExecutionId(job.jobExecutionId())
+            .withNamespace(job.namespace())
+            .withJobTemplateName("renamed_shell_job")
+            .withStartedAt(job.startedAt())
+            .withFinishedAt(job.finishedAt())
+            .withStatus(job.status())
+            .withAuditInfo(job.auditInfo())
+            .build();
+    mockListActiveJobs(renamedJob);
+    when(entityStore.delete(NameIdentifierUtil.ofJob(metalake, job.name()), 
Entity.EntityType.JOB))
+        .thenReturn(true);
+    File legacyDir = new File(testStagingDir, metalake + "/shell_job/" + 
job.name());
+    Assertions.assertTrue(legacyDir.mkdirs());
+    Assertions.assertTrue(new File(legacyDir, "output.log").createNewFile());
+    File otherJobDir = new File(testStagingDir, metalake + "/shell_job/job-" + 
(job.id() + 1));
+    Assertions.assertTrue(otherJobDir.mkdirs());
+
+    Assertions.assertDoesNotThrow(() -> jobManager.cleanUpStagingDirs());
+
+    Assertions.assertFalse(legacyDir.exists());
+    Assertions.assertTrue(otherJobDir.isDirectory());
+  }
+
+  @Test
+  public void testDeleteJobTemplateDeletesLegacyStagingDirOfRenamedMetalake() 
throws IOException {
+    // The job ran in the legacy layout under the metalake's former name.
+    JobEntity job = expiredJob();
+    
doReturn(Collections.singletonList(job)).when(jobManager).listJobs(metalake, 
Optional.of("t"));
+    doReturn(true)
+        .when(entityStore)
+        .delete(NameIdentifierUtil.ofJobTemplate(metalake, "t"), 
Entity.EntityType.JOB_TEMPLATE);
+    File legacyDir = new File(testStagingDir, "old_metalake_name/t/" + 
job.name());
+    Assertions.assertTrue(legacyDir.mkdirs());
+
+    Assertions.assertTrue(jobManager.deleteJobTemplate(metalake, "t"));
+
+    Assertions.assertFalse(legacyDir.exists());
+    // Only the job's own directory is deleted, never the template or metalake 
directory.
+    Assertions.assertTrue(legacyDir.getParentFile().isDirectory());
+  }
+
+  @Test
+  public void testDeleteJobStagingDirPrefersCurrentLayout() throws IOException 
{
+    JobEntity job = expiredJob();
+    
doReturn(Collections.singletonList(job)).when(jobManager).listJobs(metalake, 
Optional.of("t"));
+    doReturn(true)
+        .when(entityStore)
+        .delete(NameIdentifierUtil.ofJobTemplate(metalake, "t"), 
Entity.EntityType.JOB_TEMPLATE);
+    File jobStagingDir = jobManager.jobStagingDir(job.id());
+    Assertions.assertTrue(jobStagingDir.mkdirs());
+
+    Assertions.assertTrue(jobManager.deleteJobTemplate(metalake, "t"));
+
+    Assertions.assertFalse(jobStagingDir.exists());
+    verify(jobManager, never()).findLegacyJobStagingDirs(job.id());
+  }
+
+  @Test
+  public void testFindLegacyJobStagingDirs() throws IOException {
+    long jobId = idGenerator.nextId();
+    String jobDirName = JobHandle.JOB_ID_PREFIX + jobId;
+    File stagingDir = new File(testStagingDir).getAbsoluteFile();
+    File legacyDir = new File(stagingDir, metalake + "/shell_job/" + 
jobDirName);
+    Assertions.assertTrue(legacyDir.mkdirs());
+
+    // Not the job's directory: a file of the same name, a directory at the 
wrong depth, a hidden
+    // top-level directory and the directory of the current layout.
+    File sameNameFile = new File(stagingDir, metalake + "/other_job/" + 
jobDirName);
+    Assertions.assertTrue(sameNameFile.getParentFile().mkdirs());
+    Assertions.assertTrue(sameNameFile.createNewFile());
+    Assertions.assertTrue(new File(stagingDir, metalake + "/" + 
jobDirName).mkdirs());
+    Assertions.assertTrue(new File(stagingDir, ".job-output-index/x/" + 
jobDirName).mkdirs());
+    Assertions.assertTrue(jobManager.jobStagingDir(jobId).mkdirs());
+    Assertions.assertTrue(
+        new File(jobManager.jobStagingDir(jobId).getParentFile(), "x/" + 
jobDirName).mkdirs());
+    Assertions.assertTrue(new File(stagingDir, "a_file").createNewFile());
+
+    // Symbolic links are never followed, so nothing outside the staging 
directory is found.
+    File outsideDir = 
Files.createTempDirectory("gravitino-test-outside-staging").toFile();
+    try {
+      Assertions.assertTrue(new File(outsideDir, "t/" + jobDirName).mkdirs());
+      Files.createSymbolicLink(
+          new File(stagingDir, "linked_metalake").toPath(), 
outsideDir.toPath());
+      Assertions.assertTrue(new File(stagingDir, "metalake_2").mkdirs());
+      Files.createSymbolicLink(
+          new File(stagingDir, "metalake_2/linked_template").toPath(),
+          new File(outsideDir, "t").toPath());
+      File linkedJobTemplateDir = new File(stagingDir, metalake + 
"/linked_job");
+      Assertions.assertTrue(linkedJobTemplateDir.mkdirs());
+      Files.createSymbolicLink(
+          new File(linkedJobTemplateDir, jobDirName).toPath(),
+          new File(outsideDir, "t/" + jobDirName).toPath());
+    } catch (IOException e) {
+      FileUtils.deleteDirectory(outsideDir);
+      throw e;
+    }
+
+    try {
+      Assertions.assertEquals(
+          ImmutableList.of(legacyDir.getAbsolutePath()),
+          jobManager.findLegacyJobStagingDirs(jobId).stream()
+              .map(File::getAbsolutePath)
+              .collect(Collectors.toList()));
+      Assertions.assertTrue(
+          jobManager.findLegacyJobStagingDirs(idGenerator.nextId()).isEmpty(),
+          "No directory of an unknown job");
+    } finally {
+      FileUtils.deleteDirectory(outsideDir);
+    }
+  }
+
+  @Test
+  public void 
testCleanUpStagingDirsDeletesLegacyStagingDirOfNestedTemplateName()
+      throws IOException {
+    // Template names may contain '/', which the legacy layout put in the path 
as is, so the job
+    // directory is nested deeper. The template was never renamed: this is an 
upgraded job.
+    JobEntity job = expiredJob();
+    JobEntity nestedTemplateJob =
+        JobEntity.builder()
+            .withId(job.id())
+            .withJobExecutionId(job.jobExecutionId())
+            .withNamespace(job.namespace())
+            .withJobTemplateName("team/etl")
+            .withStartedAt(job.startedAt())
+            .withFinishedAt(job.finishedAt())
+            .withStatus(job.status())
+            .withAuditInfo(job.auditInfo())
+            .build();
+    mockListActiveJobs(nestedTemplateJob);
+    when(entityStore.delete(NameIdentifierUtil.ofJob(metalake, job.name()), 
Entity.EntityType.JOB))
+        .thenReturn(true);
+    File legacyDir = new File(testStagingDir, metalake + "/team/etl/" + 
job.name());
+    Assertions.assertTrue(legacyDir.mkdirs());
+
+    Assertions.assertDoesNotThrow(() -> jobManager.cleanUpStagingDirs());
+
+    Assertions.assertFalse(legacyDir.exists());
+  }
+
+  @Test
+  public void testFindLegacyJobStagingDirsOfNestedTemplateNames() throws 
IOException {
+    long jobId = idGenerator.nextId();
+    String jobDirName = JobHandle.JOB_ID_PREFIX + jobId;
+    File stagingDir = new File(testStagingDir).getAbsoluteFile();
+    File nestedDir = new File(stagingDir, metalake + "/team/etl/" + 
jobDirName);
+    Assertions.assertTrue(nestedDir.mkdirs());
+    File deeplyNestedDir = new File(stagingDir, metalake + "/a/b/c/d/" + 
jobDirName);
+    Assertions.assertTrue(deeplyNestedDir.mkdirs());
+
+    // The staging directory of another job is not descended into: a directory 
of its own files
+    // named like this job is not this job's staging directory.
+    File otherJobDir =
+        new File(
+            stagingDir,
+            metalake
+                + "/shell_job/"
+                + JobHandle.JOB_ID_PREFIX
+                + (jobId + 1)
+                + File.separator
+                + jobDirName);
+    Assertions.assertTrue(otherJobDir.mkdirs());
+
+    // A template name nested deeper than the lookup goes keeps its directory.
+    StringBuilder tooDeepPath = new StringBuilder(metalake);
+    for (int i = 0; i < 20; i++) {
+      tooDeepPath.append(File.separator).append("d").append(i);
+    }
+    File tooDeepDir = new File(stagingDir, tooDeepPath + File.separator + 
jobDirName);
+    Assertions.assertTrue(tooDeepDir.mkdirs());
+
+    Assertions.assertEquals(
+        ImmutableSet.of(nestedDir.getAbsolutePath(), 
deeplyNestedDir.getAbsolutePath()),
+        jobManager.findLegacyJobStagingDirs(jobId).stream()
+            .map(File::getAbsolutePath)
+            .collect(Collectors.toSet()));
+  }
+
+  @Test
+  public void testFindLegacyJobStagingDirsSkipsUnreadableDirectory() throws 
IOException {
+    long jobId = idGenerator.nextId();
+    File stagingDir = new File(testStagingDir);
+    File legacyDir =
+        new File(stagingDir, metalake + "/shell_job/" + 
JobHandle.JOB_ID_PREFIX + jobId);
+    Assertions.assertTrue(legacyDir.mkdirs());
+    // E.g. "lost+found" when the staging directory is the root of a file 
system.
+    File unreadableDir = new File(stagingDir, "lost+found");
+    Assertions.assertTrue(unreadableDir.mkdirs());
+    Assertions.assertTrue(unreadableDir.setReadable(false, false));
+
+    try {
+      Assertions.assertEquals(
+          ImmutableList.of(legacyDir.getAbsolutePath()),
+          jobManager.findLegacyJobStagingDirs(jobId).stream()
+              .map(File::getAbsolutePath)
+              .collect(Collectors.toList()));
+    } finally {
+      Assertions.assertTrue(unreadableDir.setReadable(true, false));
+    }
+  }
+
   @Test
   public void testCleanUpStagingDirs() throws IOException, 
InterruptedException {
     JobEntity job = newJobEntity("shell_job", JobHandle.Status.STARTED);
diff --git 
a/core/src/test/java/org/apache/gravitino/job/TestJobManagerMultiNode.java 
b/core/src/test/java/org/apache/gravitino/job/TestJobManagerMultiNode.java
index 9f830bdb4c..f2f50cda2d 100644
--- a/core/src/test/java/org/apache/gravitino/job/TestJobManagerMultiNode.java
+++ b/core/src/test/java/org/apache/gravitino/job/TestJobManagerMultiNode.java
@@ -24,10 +24,14 @@ import com.google.common.collect.Lists;
 import java.io.File;
 import java.io.IOException;
 import java.nio.file.Files;
+import java.nio.file.Path;
 import java.time.Instant;
 import java.util.Collections;
 import java.util.EnumSet;
+import java.util.List;
 import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
+import java.util.stream.Stream;
 import org.apache.commons.io.FileUtils;
 import org.apache.commons.lang3.reflect.FieldUtils;
 import org.apache.gravitino.Config;
@@ -42,6 +46,7 @@ import org.apache.gravitino.job.local.LocalJobExecutor;
 import org.apache.gravitino.job.local.LocalJobExecutorConfigs;
 import org.apache.gravitino.lock.LockManager;
 import org.apache.gravitino.meta.AuditInfo;
+import org.apache.gravitino.meta.BaseMetalake;
 import org.apache.gravitino.meta.JobEntity;
 import org.apache.gravitino.meta.JobTemplateEntity;
 import org.apache.gravitino.storage.RandomIdGenerator;
@@ -227,7 +232,7 @@ public class TestJobManagerMultiNode extends 
TestJDBCBackend {
     String newName = ECHO_TEMPLATE + "_renamed";
     nodeA.alterJobTemplate(METALAKE, ECHO_TEMPLATE, 
JobTemplateChange.rename(newName));
 
-    // The job reports the new template name, while its staging directory 
keeps the old one.
+    // The job reports the new template name, which its staging directory 
doesn't depend on.
     JobEntity jobWithOutput = nodeB.getJob(METALAKE, job.name(), true);
     Assertions.assertEquals(newName, jobWithOutput.jobTemplateName());
     Assertions.assertEquals(ImmutableList.of("hello b"), 
jobWithOutput.stdout());
@@ -235,7 +240,7 @@ public class TestJobManagerMultiNode extends 
TestJDBCBackend {
 
   @TestTemplate
   public void testGetJobOutputOfTemplateNamedWithSpecialCharacters() throws 
IOException {
-    // Template names are not restricted, and name a directory level of the 
job staging directory.
+    // Template names are not restricted, and must not affect running the job 
or reading its output.
     for (String templateName : ImmutableList.of("etl job \"v2\" 中文 #1", 
"team/etl")) {
       backend.insert(
           newScriptJobTemplateEntity(
@@ -300,6 +305,87 @@ public class TestJobManagerMultiNode extends 
TestJDBCBackend {
     return job;
   }
 
+  @TestTemplate
+  public void testStagingDirIsCleanedUpAfterTemplateRenamed() throws 
IOException {
+    JobEntity job = runFinishedJobOnNodeA();
+    Assertions.assertTrue(jobStagingDir(job).isDirectory());
+
+    // Before the fix, the cleanup rebuilt the staging path from the new 
template name, missed the
+    // directory and deleted only the job entity, leaking the directory.
+    nodeA.alterJobTemplate(METALAKE, TEMPLATE, 
JobTemplateChange.rename("renamed_sleep_job"));
+    Assertions.assertEquals("renamed_sleep_job", 
getJob(job.name()).jobTemplateName());
+    moveJobTimestampsBack(job.name());
+    nodeB.cleanUpStagingDirs();
+
+    Assertions.assertFalse(jobExists(job.name()));
+    assertNoStagingDirLeft(job);
+  }
+
+  @TestTemplate
+  public void testStagingDirIsDeletedWithRenamedTemplate() throws IOException {
+    JobEntity job = runFinishedJobOnNodeA();
+
+    nodeA.alterJobTemplate(METALAKE, TEMPLATE, 
JobTemplateChange.rename("renamed_sleep_job"));
+    nodeA.alterJobTemplate(
+        METALAKE, "renamed_sleep_job", 
JobTemplateChange.rename("renamed_twice_sleep_job"));
+    Assertions.assertTrue(nodeB.deleteJobTemplate(METALAKE, 
"renamed_twice_sleep_job"));
+
+    Assertions.assertFalse(jobExists(job.name()));
+    assertNoStagingDirLeft(job);
+  }
+
+  @TestTemplate
+  public void testStagingDirIsCleanedUpAfterMetalakeRenamed() throws 
IOException {
+    JobEntity job = runFinishedJobOnNodeA();
+    moveJobTimestampsBack(job.name());
+
+    String newMetalake = METALAKE + "_renamed";
+    entityStore.update(
+        NameIdentifierUtil.ofMetalake(METALAKE),
+        BaseMetalake.class,
+        Entity.EntityType.METALAKE,
+        metalake ->
+            BaseMetalake.builder()
+                .withId(metalake.id())
+                .withName(newMetalake)
+                .withComment(metalake.comment())
+                .withProperties(metalake.properties())
+                .withAuditInfo(metalake.auditInfo())
+                .withVersion(metalake.getVersion())
+                .build());
+    nodeB.cleanUpStagingDirs();
+
+    Assertions.assertThrows(
+        NoSuchJobException.class, () -> nodeB.getJob(newMetalake, job.name(), 
false));
+    assertNoStagingDirLeft(job);
+  }
+
+  private JobEntity runFinishedJobOnNodeA() throws IOException {
+    JobEntity job = nodeA.runJob(METALAKE, TEMPLATE, 
ImmutableMap.of("seconds", "0"));
+    Awaitility.await()
+        .atMost(1, TimeUnit.MINUTES)
+        .until(() -> executorA.getJobStatus(job.jobExecutionId()) == 
JobHandle.Status.SUCCEEDED);
+    nodeA.pullAndUpdateJobStatus();
+    Assertions.assertEquals(JobHandle.Status.SUCCEEDED, 
getJob(job.name()).status());
+    return job;
+  }
+
+  private File jobStagingDir(JobEntity job) {
+    return new File(testDir, "staging/job-runs/" + job.name());
+  }
+
+  // Checks the whole staging directory rather than the expected path, so a 
directory left in any
+  // layout is caught.
+  private void assertNoStagingDirLeft(JobEntity job) throws IOException {
+    try (Stream<Path> paths = Files.walk(new File(testDir, 
"staging").toPath())) {
+      List<Path> left =
+          paths
+              .filter(path -> path.getFileName().toString().equals(job.name()))
+              .collect(Collectors.toList());
+      Assertions.assertTrue(left.isEmpty(), "Staging directory left behind: " 
+ left);
+    }
+  }
+
   private JobEntity runLongJobOnNodeA() throws IOException {
     JobEntity job = nodeA.runJob(METALAKE, TEMPLATE, 
ImmutableMap.of("seconds", "600"));
     Awaitility.await()
diff --git 
a/core/src/test/java/org/apache/gravitino/job/local/TestLocalJobExecutor.java 
b/core/src/test/java/org/apache/gravitino/job/local/TestLocalJobExecutor.java
index 6c5441dd79..ee2d0dae25 100644
--- 
a/core/src/test/java/org/apache/gravitino/job/local/TestLocalJobExecutor.java
+++ 
b/core/src/test/java/org/apache/gravitino/job/local/TestLocalJobExecutor.java
@@ -785,22 +785,23 @@ public class TestLocalJobExecutor {
 
   @Test
   public void testSubmitJobWritesOutputIndexRelativeToStagingDir() throws 
IOException {
-    // Laid out like the {metalake}/{template}/job-{id} staging directory of 
JobManager.
-    File jobDir = new File(workingDir, "metalake/template/job-1");
+    // Laid out like the job-runs/job-{id} staging directory of JobManager.
+    File jobDir = new File(workingDir, "job-runs/job-1");
     Assertions.assertTrue(jobDir.mkdirs());
     String jobId = runSucceededJob(jobDir);
 
     JsonNode index = 
JsonUtils.anyFieldMapper().readTree(outputIndexFile(jobId));
     Assertions.assertEquals(1, index.get("version").intValue());
     Assertions.assertEquals(
-        workingDir.getName() + "/metalake/template/job-1", 
index.get("workingDir").textValue());
+        workingDir.getName() + "/job-runs/job-1", 
index.get("workingDir").textValue());
   }
 
   @Test
   public void testOutputIndexKeepsSpecialCharactersInWorkingDir() throws 
IOException {
-    // Job template names are not restricted, so the staging directory may 
contain any character a
-    // file name can. They must survive the JSON encoding and the '/'-joining 
unchanged. Some of
-    // these characters are only valid in POSIX file names, like the shell job 
itself.
+    // The working directory may contain any character a file name can, e.g. a 
staging directory of
+    // an earlier version, which was named after the job template. It must 
survive the JSON encoding
+    // and the '/'-joining unchanged. Some of these characters are only valid 
in POSIX file names,
+    // like the shell job itself.
     String specialName =
         "a b \"quoted\" back\\slash 中文 \t tab \n newline %20 
#!$&'()*+,;=@[]{}~`^|<>?";
     File jobDir =
diff --git a/docs/table-maintenance-service/optimizer-cli-reference.md 
b/docs/table-maintenance-service/optimizer-cli-reference.md
index 81c35f0a1a..c93b2bde9b 100644
--- a/docs/table-maintenance-service/optimizer-cli-reference.md
+++ b/docs/table-maintenance-service/optimizer-cli-reference.md
@@ -347,7 +347,7 @@ CALL `rest_catalog`.system.expire_snapshots(
 
 ```bash
 curl -sS "http://localhost:8090/api/metalakes/test/jobs/{job_id}"; | jq 
'.job.state'
-cat 
/tmp/gravitino/jobs/staging/test/builtin-iceberg-expire-snapshots/{job_id}/stdout.log
+cat /tmp/gravitino/jobs/staging/job-runs/{job_id}/output.log
 ```
 
 A successful run reports its state as `SUCCEEDED` and logs the counts it 
removed:
diff --git a/docs/table-maintenance-service/optimizer-troubleshooting.md 
b/docs/table-maintenance-service/optimizer-troubleshooting.md
index 099b9b4814..98e738021b 100644
--- a/docs/table-maintenance-service/optimizer-troubleshooting.md
+++ b/docs/table-maintenance-service/optimizer-troubleshooting.md
@@ -12,7 +12,7 @@ license: "This software is licensed under the Apache License 
version 2."
 
 Failures fall into three groups, matching where they occur in the workflow. 
Command and argument errors surface immediately. Evaluation problems produce no 
output rather than an error, which is what makes them confusing. Execution 
failures happen inside Spark, so the real message is in the staging log rather 
than the API response.
 
-Staging logs live under 
`/tmp/gravitino/jobs/staging/{metalake}/{job_template_name}/{job_id}/`, 
controlled by `gravitino.job.stagingDir`. Read `error.log` for failures and 
`output.log` for results.
+Staging logs live under `/tmp/gravitino/jobs/staging/job-runs/{job_id}/`, 
controlled by `gravitino.job.stagingDir`. Read `error.log` for failures and 
`output.log` for results.
 
 ## Command and Argument Errors
 

Reply via email to