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


##########
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:
   Good catch, fixed in aead588a7: output is only read from a working directory 
whose real path equals the real staging root plus the indexed relative path (so 
a symlinked working directory, or a symlinked directory above it such as one 
pointing at another metalake's job, is rejected), and output.log/error.log are 
opened with NOFOLLOW_LINKS. Added 
TestLocalJobExecutor#testGetJobOutputDoesNotFollowSymlinks (file, working dir 
and parent dir symlinks, read from a second executor) and 
TestJobManagerMultiNode#testGetJobOutputDoesNotFollowSymlinkFromAnotherNode, 
where the job replaces its output.log with an absolute symlink and node B reads 
it.
   
   _🤖 Addressed by [Claude Code](https://claude.com/claude-code)_



##########
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:
   Done in aead588a7: the staging directory is now injected when the 
instantiated executor is a LocalJobExecutor (instanceof), and 
TestJobExecutorFactory#testLocalJobExecutorSubclassUsesJobStagingDir covers a 
subclass that inherits initialize().
   
   _🤖 Addressed by [Claude Code](https://claude.com/claude-code)_



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