yuqi1129 commented on code in PR #13373:
URL: https://github.com/apache/gravitino/pull/13373#discussion_r4070107522


##########
core/src/main/java/org/apache/gravitino/job/local/LocalJobExecutor.java:
##########
@@ -404,29 +505,201 @@ void cleanupJobStatus() {
       jobStatus
           .entrySet()
           .removeIf(
-              entry -> {
-                boolean expired =
-                    entry.getValue().getRight() != UNEXPIRED_TIME_IN_MS
-                        && (currentTime - entry.getValue().getRight()) >= 
jobStatusKeepTimeInMs;
-                if (expired) {
-                  jobWorkingDirs.remove(entry.getKey());
-                }
-                return expired;
-              });
+              entry ->
+                  entry.getValue().getRight() != UNEXPIRED_TIME_IN_MS
+                      && (currentTime - entry.getValue().getRight()) >= 
jobStatusKeepTimeInMs);
+    }
+  }
+
+  /**
+   * Removes the output index files whose job staging directory no longer 
exists. The staging
+   * directories themselves are removed by JobManager, so this never deletes 
any job output.
+   */
+  @VisibleForTesting
+  void cleanupOutputIndexes() {
+    long now = System.currentTimeMillis();
+    // An exception escaping a scheduled task would cancel all its later runs.
+    try (DirectoryStream<Path> indexFiles =
+        Files.newDirectoryStream(
+            outputIndexDir, LOCAL_JOB_PREFIX + "*" + 
OUTPUT_INDEX_FILE_SUFFIX)) {
+      for (Path indexFile : indexFiles) {
+        // close() interrupts this thread, and every later read would then 
fail as well.
+        if (Thread.currentThread().isInterrupted()) {
+          return;
+        }
+        try {
+          // Anything but an index file, e.g. a directory of a colliding 
legacy metalake name.
+          if (!Files.isRegularFile(indexFile, LinkOption.NOFOLLOW_LINKS)) {
+            continue;
+          }
+          if (now - Files.getLastModifiedTime(indexFile).toMillis() < 
OUTPUT_INDEX_MIN_AGE_IN_MS) {
+            continue;
+          }
+          Path workingDir = parseOutputIndex(Files.readAllBytes(indexFile));
+          // An index this server can't interpret may come from a newer 
server, so keep it. Only
+          // a confirmed missing directory counts: notExists() is false on I/O 
errors too, so a
+          // storage hiccup never removes an index.
+          if (workingDir != null && Files.notExists(workingDir)) {
+            Files.deleteIfExists(indexFile);
+            LOG.debug("Removed output index {} of a deleted job staging 
directory", indexFile);
+          }
+        } catch (ClosedByInterruptException e) {
+          return;
+        } catch (NoSuchFileException e) {
+          // Removed concurrently, e.g. by another server sharing the staging 
directory.
+        } catch (IOException | RuntimeException e) {
+          LOG.warn("Failed to clean up job output index {}", indexFile, e);
+        }
+      }
+    } catch (IOException | RuntimeException e) {
+      LOG.warn("Failed to clean up job output indexes under {}", 
outputIndexDir, e);
     }
   }
 
   private List<String> getJobOutput(String jobId, String fileName, int 
maxLines, int maxBytes) {
-    File workingDir = getWorkingDir(jobId);
+    Path workingDir = locateWorkingDir(jobId);
     if (workingDir == null) {
       return ImmutableList.of();
     }
-    return readLastLines(new File(workingDir, fileName), maxLines, maxBytes);
+    return readLastLines(workingDir.resolve(fileName).toFile(), maxLines, 
maxBytes);

Review Comment:
   The containment check in `parseOutputIndex` is lexical, but `readLastLines` 
follows symlinks. A local shell job can replace its `output.log` with an 
absolute symlink after the process starts; when another server sharing the 
staging directory handles `getJob(includeOutput=true)`, it resolves that 
symlink against its own filesystem and can return a file outside staging. I 
reproduced the path behavior with Java (`startsWith(stagingRoot)` true, read 
from an outside file). Could we reject symlinks or validate the real output 
path against the real staging root before reading, and add a cross-node symlink 
test?



##########
core/src/main/java/org/apache/gravitino/job/JobExecutorFactory.java:
##########
@@ -60,7 +61,13 @@ public static JobExecutor create(Config config) {
         jobExecutorName);
 
     Map<String, String> configs =
-        config.getConfigsWithPrefix(JOB_EXECUTOR_CONF_PREFIX + jobExecutorName 
+ ".");
+        Maps.newHashMap(
+            config.getConfigsWithPrefix(JOB_EXECUTOR_CONF_PREFIX + 
jobExecutorName + "."));
+    if (LocalJobExecutor.class.getCanonicalName().equals(clzName)) {

Review Comment:
   This exact class-name check misses a configured subclass of 
`LocalJobExecutor` that inherits `initialize()`. Such an executor no longer 
receives `stagingDir`, so the new non-blank check fails during server startup. 
Could we inject the staging directory based on the instantiated executor type 
(for example, `instanceof LocalJobExecutor`) and cover a subclass in 
`TestJobExecutorFactory`?



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