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

jerryshao 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 43432b802f [#13365] improvement(core): Record job start and finish 
times reported by the job executor (#13381)
43432b802f is described below

commit 43432b802fc3e92e532c878dbcc8d636094c2e4a
Author: Jerry Shao <[email protected]>
AuthorDate: Thu Sep 24 17:37:48 2026 +0800

    [#13365] improvement(core): Record job start and finish times reported by 
the job executor (#13381)
    
    ### What changes were proposed in this pull request?
    
    The job executor now reports when a job actually started and finished,
    instead of Gravitino inferring both times from its status poll.
    
    - **SPI**: new `JobExecutor#getJobExecutionInfo`, returning a
    `JobExecutionInfo` snapshot: status, `startedAt`, `finishedAt`.
    - The timestamps are attributes of the job, not of a status. Once a job
    has started, every later snapshot carries its start time, including the
    terminal one. So nothing is lost when a job goes from `QUEUED` to
    `SUCCEEDED` between two polls.
    - It is abstract, so an executor that doesn't implement it fails to
    compile, instead of silently reporting no times.
      - `getJobStatus` becomes a default shortcut of `getJobExecutionInfo`.
    - **`LocalJobExecutor`** records each job's real start and finish times,
    and keeps a `JobExecutionInfo` per job instead of a `Pair<Status,
    Long>`.
    - **`JobManager`**:
    - polls with `getJobExecutionInfo`, and uses the reported times when
    present. The reported start time also replaces one recorded by an
    earlier poll;
    - without reported times, keeps the current behavior: poll time for an
    observed `STARTED`, no `startedAt` for a job never observed running;
    - corrects inconsistent reported data rather than failing the poll: a
    time on the wrong status, one that doesn't fit in epoch milliseconds, or
    one more than a day ahead is dropped; a time earlier than `queuedAt`
    (clock skew) is raised to it; and a `startedAt` later than `finishedAt`
    is dropped. As the same report is corrected on every poll, the
    corrections are logged at debug level;
    - catches a failure on one job, so it no longer escapes the scheduled
    status pull and cancels all its later runs;
    - updates a job when its status or timestamps change, and skips the
    write when nothing changes;
    - takes `queuedAt` before submitting the job, so it is never later than
    the reported start time.
    - **Docs**: the timestamp semantics in `manage-jobs-in-gravitino.md`,
    the `getJobExecutionInfo` requirement in `custom-job-executor.md`, and
    the meaning of a null `startedAt` on a finished job in `JobHandle`,
    `JobInfo`, OpenAPI and the Python client.
    
    ### Why are the changes needed?
    
    `startedAt` was only set when a status poll happened to observe the job
    in `STARTED`. The poll runs every `gravitino.job.statusPullIntervalInMs`
    (5 minutes by default), so any job shorter than that usually finished
    with no `startedAt`. `finishedAt` was the time of the poll that first
    saw the job finished, so it could be up to one interval late. Durations
    therefore reflected the polling schedule rather than the job, and queue
    time couldn't be told apart from run time.
    
    Fix: #13365
    
    ### Does this PR introduce _any_ user-facing change?
    
    - `startedAt` and `finishedAt` are now the actual times for jobs run by
    the local job executor, including jobs that start and finish between two
    polls.
    - **Job executor SPI (breaking)**: custom job executors must implement
    `getJobExecutionInfo`, or they fail to compile. Overriding only
    `getJobStatus` is no longer enough.
    - No REST API or configuration change.
    
    ### How was this patch tested?
    
    - New `TestJobExecutionInfo`: the builder, the `started`/`finished`
    helpers, and `getJobStatus` derived from `getJobExecutionInfo`.
    - `TestJobManager`:
      - reported times are used;
      - a reported start time replaces the poll time;
    - an update with only a timestamp change is written, and one with no
    change is skipped;
    - inconsistent, clock-skewed and unusable reported times are corrected,
    including a poll-recorded start time later than the reported finish
    time;
      - a CANCELLING job is only updated once it finishes;
      - a failure on one job doesn't stop the status pull of the others;
      - `queuedAt` is taken before the submission.
    - `TestLocalJobExecutor`: the timestamps for succeeded, failed,
    cancelled-while-queued and cancelled-while-running jobs, and the cleanup
    of finished jobs.
    - `TestJobManagerMultiNode`: a job that finishes before it is polled
    records its actual start and finish times.
    - `./gradlew :core:test :server:test -PskipITs`, `spotlessCheck`,
    `:docs:build`, `:api:javadoc`, `:core:javadoc`.
    
    🤖 Generated with [Claude Code](https://claude.com/claude-code)
    
    ---------
    
    Co-authored-by: Claude Opus 5 <[email protected]>
---
 .../java/org/apache/gravitino/job/JobHandle.java   |   6 +-
 .../client-python/gravitino/api/job/job_handle.py  |   3 +-
 .../gravitino/connector/job/JobExecutionInfo.java  | 113 ++++++
 .../gravitino/connector/job/JobExecutor.java       |  62 +++-
 .../apache/gravitino/job/JobExecutorFactory.java   |  39 +-
 .../java/org/apache/gravitino/job/JobManager.java  | 369 +++++++++++++------
 .../gravitino/job/local/LocalJobExecutor.java      |  87 ++---
 .../gravitino/listener/api/info/JobInfo.java       |   6 +-
 .../connector/job/TestJobExecutionInfo.java        | 139 +++++++
 .../gravitino/job/TestJobExecutorFactory.java      | 108 +++++-
 .../org/apache/gravitino/job/TestJobManager.java   | 409 +++++++++++++++++++--
 .../gravitino/job/TestJobManagerMultiNode.java     |  23 ++
 .../gravitino/job/local/TestLocalJobExecutor.java  | 107 +++++-
 docs/development/custom-job-executor.md            |   7 +
 docs/manage-jobs-in-gravitino.md                   |  22 ++
 docs/open-api/jobs.yaml                            |   2 +-
 16 files changed, 1303 insertions(+), 199 deletions(-)

diff --git a/api/src/main/java/org/apache/gravitino/job/JobHandle.java 
b/api/src/main/java/org/apache/gravitino/job/JobHandle.java
index d9463fb053..b6fc536aaf 100644
--- a/api/src/main/java/org/apache/gravitino/job/JobHandle.java
+++ b/api/src/main/java/org/apache/gravitino/job/JobHandle.java
@@ -88,7 +88,11 @@ public interface JobHandle {
   /**
    * Get the time when the job started execution.
    *
-   * @return the started time of the job, or null if the job has not started 
execution yet
+   * <p>A finished job may also have no started time, if the job executor 
doesn't report when the
+   * job started and Gravitino didn't observe the job running.
+   *
+   * @return the started time of the job, or null if the job has not started 
execution yet, or the
+   *     started time is unknown
    */
   @Nullable
   default Instant startedAt() {
diff --git a/clients/client-python/gravitino/api/job/job_handle.py 
b/clients/client-python/gravitino/api/job/job_handle.py
index 98dbda8d44..39473e6c4b 100644
--- a/clients/client-python/gravitino/api/job/job_handle.py
+++ b/clients/client-python/gravitino/api/job/job_handle.py
@@ -63,7 +63,8 @@ class JobHandle(ABC):
 
     def started_at(self) -> Optional[datetime]:
         """Returns the time the job started execution, or ``None`` if the job 
has not started
-        execution yet.
+        execution yet. A finished job may also have no started time, if the 
job executor doesn't
+        report when the job started and Gravitino didn't observe the job 
running.
         """
         raise NotImplementedError("started_at is not implemented")
 
diff --git 
a/core/src/main/java/org/apache/gravitino/connector/job/JobExecutionInfo.java 
b/core/src/main/java/org/apache/gravitino/connector/job/JobExecutionInfo.java
new file mode 100644
index 0000000000..b433eb6320
--- /dev/null
+++ 
b/core/src/main/java/org/apache/gravitino/connector/job/JobExecutionInfo.java
@@ -0,0 +1,113 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *  http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.gravitino.connector.job;
+
+import com.google.common.base.Preconditions;
+import java.time.Instant;
+import javax.annotation.Nullable;
+import lombok.AccessLevel;
+import lombok.AllArgsConstructor;
+import lombok.Builder;
+import lombok.EqualsAndHashCode;
+import lombok.Getter;
+import lombok.NonNull;
+import lombok.ToString;
+import lombok.experimental.Accessors;
+import org.apache.gravitino.annotation.DeveloperApi;
+import org.apache.gravitino.job.JobHandle;
+
+/**
+ * A snapshot of a job's execution state, as reported by a {@link JobExecutor}.
+ *
+ * <p>The timestamps are attributes of the job, not of a particular status. 
Once a job has started,
+ * every later snapshot carries its start time, including the terminal ones, 
so no information is
+ * lost when Gravitino doesn't observe every status the job goes through.
+ *
+ * <p>More optional fields may be added in the future. Create instances with 
{@link
+ * #of(JobHandle.Status)} or {@code builder()}, so that the existing code 
keeps working when fields
+ * are added.
+ */
+@DeveloperApi
+@Getter
+@Accessors(fluent = true)
+@EqualsAndHashCode
+@ToString
+@AllArgsConstructor(access = AccessLevel.PRIVATE)
+@Builder(setterPrefix = "with", toBuilder = true)
+public final class JobExecutionInfo {
+
+  /** The status of the job. */
+  @NonNull private final JobHandle.Status status;
+
+  /**
+   * The time when the job actually started executing, or null if the job 
hasn't started executing,
+   * or the started time is unknown to the job executor.
+   */
+  @Nullable private final Instant startedAt;
+
+  /**
+   * The time when the job actually finished, or null if the job hasn't 
finished, or the finished
+   * time is unknown to the job executor.
+   */
+  @Nullable private final Instant finishedAt;
+
+  /**
+   * Create a job execution info that only carries the status of the job, 
without any timestamp.
+   *
+   * @param status The status of the job.
+   * @return the job execution info.
+   */
+  public static JobExecutionInfo of(JobHandle.Status status) {
+    return builder().withStatus(status).build();
+  }
+
+  /**
+   * Create a copy of this job execution info with the job moved to {@link
+   * JobHandle.Status#STARTED}.
+   *
+   * @param startedAt The time when the job started executing.
+   * @return the job execution info of the started job.
+   */
+  public JobExecutionInfo started(Instant startedAt) {
+    Preconditions.checkArgument(startedAt != null, "The started time must not 
be null");
+    return 
toBuilder().withStatus(JobHandle.Status.STARTED).withStartedAt(startedAt).build();
+  }
+
+  /**
+   * Create a copy of this job execution info with the job moved to the given 
terminal status. The
+   * started time, if any, is carried forward.
+   *
+   * @param terminalStatus The terminal status of the job, one of {@link
+   *     JobHandle.Status#SUCCEEDED}, {@link JobHandle.Status#FAILED} and 
{@link
+   *     JobHandle.Status#CANCELLED}.
+   * @param finishedAt The time when the job finished.
+   * @return the job execution info of the finished job.
+   */
+  public JobExecutionInfo finished(JobHandle.Status terminalStatus, Instant 
finishedAt) {
+    Preconditions.checkArgument(
+        terminalStatus == JobHandle.Status.SUCCEEDED
+            || terminalStatus == JobHandle.Status.FAILED
+            || terminalStatus == JobHandle.Status.CANCELLED,
+        "The status of a finished job must be SUCCEEDED, FAILED or CANCELLED, 
but got: %s",
+        terminalStatus);
+    Preconditions.checkArgument(finishedAt != null, "The finished time must 
not be null");
+    return 
toBuilder().withStatus(terminalStatus).withFinishedAt(finishedAt).build();
+  }
+}
diff --git 
a/core/src/main/java/org/apache/gravitino/connector/job/JobExecutor.java 
b/core/src/main/java/org/apache/gravitino/connector/job/JobExecutor.java
index f6e5243213..88dcbd4a57 100644
--- a/core/src/main/java/org/apache/gravitino/connector/job/JobExecutor.java
+++ b/core/src/main/java/org/apache/gravitino/connector/job/JobExecutor.java
@@ -60,15 +60,45 @@ public interface JobExecutor extends Closeable {
   String submitJob(JobTemplate jobTemplate);
 
   /**
-   * Get the status of a job by its unique identifier. The status should be 
one of the values in
-   * {@link JobHandle.Status}. The implementors should query the external job 
runner to get the job
-   * status, and map the status to the values in {@link JobHandle.Status}.
+   * Get a snapshot of the job's execution state by its unique identifier, 
including its status and
+   * when the job actually started and finished. The implementors should query 
the external job
+   * runner to get the job state, and map the status to the values in {@link 
JobHandle.Status}.
+   *
+   * <p>Gravitino pulls the job state periodically, so a job may go through 
several statuses between
+   * two pulls, for example, from {@link JobHandle.Status#QUEUED} straight to 
{@link
+   * JobHandle.Status#SUCCEEDED}. The timestamps are attributes of the job, 
not of a particular
+   * status: once the job has started, the started time must be reported in 
every later snapshot,
+   * including the terminal ones. Once reported, a timestamp should not change.
+   *
+   * <ul>
+   *   <li>The started time is null if the job hasn't started executing, or if 
it is unknown to the
+   *       job executor.
+   *   <li>The finished time is only set for a terminal status, and is null if 
it is unknown to the
+   *       job executor.
+   * </ul>
+   *
+   * <p>If the job runner can't tell when a job started or finished, leave the 
time null. Gravitino
+   * then falls back to the time it observes the job running or finished, 
which can be off by up to
+   * the job status pull interval. A job that starts and finishes between two 
pulls is never
+   * observed running, so it has no started time in that case.
+   *
+   * @param jobId The unique identifier of the job.
+   * @return The execution snapshot of the job.
+   * @throws NoSuchJobException If the job with the given identifier does not 
exist.
+   */
+  JobExecutionInfo getJobExecutionInfo(String jobId) throws NoSuchJobException;
+
+  /**
+   * Get the status of a job by its unique identifier. It is a shortcut of 
{@link
+   * #getJobExecutionInfo(String)} that only returns the status.
    *
    * @param jobId The unique identifier of the job.
    * @return The status of the job.
    * @throws NoSuchJobException If the job with the given identifier does not 
exist.
    */
-  JobHandle.Status getJobStatus(String jobId) throws NoSuchJobException;
+  default JobHandle.Status getJobStatus(String jobId) throws 
NoSuchJobException {
+    return getJobExecutionInfo(jobId).status();
+  }
 
   /**
    * Cancel a job by its unique identifier. The job runner should stop the job 
if it is currently
@@ -119,12 +149,12 @@ public interface JobExecutor extends Closeable {
    * Get the captured standard output of the job, as a list of lines.
    *
    * <p>The default implementation returns an empty list, so implementors that 
don't support output
-   * retrieval don't need to override this method. Unlike {@link 
#getJobStatus(String)}/{@link
-   * #cancelJob(String)}, this method never throws for a job the executor 
doesn't (or no longer)
-   * know about - the job entity itself may still exist even after the 
executor's own bookkeeping
-   * for its output has expired or been lost (e.g. when it's only kept in 
memory and the server
-   * restarts, or kept on storage this server can't reach), so "unknown to 
this executor" is
-   * reported as empty output, not as an error.
+   * retrieval don't need to override this method. Unlike {@link
+   * #getJobExecutionInfo(String)}/{@link #cancelJob(String)}, this method 
never throws for a job
+   * the executor doesn't (or no longer) know about - the job entity itself 
may still exist even
+   * after the executor's own bookkeeping for its output has expired or been 
lost (e.g. when it's
+   * only kept in memory and the server restarts, or kept on storage this 
server can't reach), so
+   * "unknown to this executor" is reported as empty output, not as an error.
    *
    * @param jobId The unique identifier of the job.
    * @param maxLines The maximum number of (most recent) lines to return, 
resolved by the caller
@@ -142,12 +172,12 @@ public interface JobExecutor extends Closeable {
    * Get the captured standard error output of the job, as a list of lines.
    *
    * <p>The default implementation returns an empty list, so implementors that 
don't support output
-   * retrieval don't need to override this method. Unlike {@link 
#getJobStatus(String)}/{@link
-   * #cancelJob(String)}, this method never throws for a job the executor 
doesn't (or no longer)
-   * know about - the job entity itself may still exist even after the 
executor's own bookkeeping
-   * for its output has expired or been lost (e.g. when it's only kept in 
memory and the server
-   * restarts, or kept on storage this server can't reach), so "unknown to 
this executor" is
-   * reported as empty output, not as an error.
+   * retrieval don't need to override this method. Unlike {@link
+   * #getJobExecutionInfo(String)}/{@link #cancelJob(String)}, this method 
never throws for a job
+   * the executor doesn't (or no longer) know about - the job entity itself 
may still exist even
+   * after the executor's own bookkeeping for its output has expired or been 
lost (e.g. when it's
+   * only kept in memory and the server restarts, or kept on storage this 
server can't reach), so
+   * "unknown to this executor" is reported as empty output, not as an error.
    *
    * @param jobId The unique identifier of the job.
    * @param maxLines The maximum number of (most recent) lines to return, 
resolved by the caller
diff --git 
a/core/src/main/java/org/apache/gravitino/job/JobExecutorFactory.java 
b/core/src/main/java/org/apache/gravitino/job/JobExecutorFactory.java
index dc0ec158a6..46b5c35554 100644
--- a/core/src/main/java/org/apache/gravitino/job/JobExecutorFactory.java
+++ b/core/src/main/java/org/apache/gravitino/job/JobExecutorFactory.java
@@ -18,9 +18,12 @@
  */
 package org.apache.gravitino.job;
 
+import com.google.common.annotations.VisibleForTesting;
 import com.google.common.base.Preconditions;
 import com.google.common.collect.ImmutableMap;
 import com.google.common.collect.Maps;
+import java.lang.reflect.Method;
+import java.lang.reflect.Modifier;
 import java.util.Map;
 import org.apache.commons.lang3.StringUtils;
 import org.apache.gravitino.Config;
@@ -64,8 +67,10 @@ public class JobExecutorFactory {
         Maps.newHashMap(
             config.getConfigsWithPrefix(JOB_EXECUTOR_CONF_PREFIX + 
jobExecutorName + "."));
     try {
+      Class<?> jobExecutorClass = Class.forName(clzName);
+      checkJobExecutorClass(jobExecutorClass);
       JobExecutor jobExecutor =
-          (JobExecutor) 
Class.forName(clzName).getDeclaredConstructor().newInstance();
+          (JobExecutor) 
jobExecutorClass.getDeclaredConstructor().newInstance();
       if (jobExecutor instanceof LocalJobExecutor) {
         // The local job executor, and any subclass of it, keeps its output 
index under the job
         // staging directory, so it must resolve paths against exactly the 
directory JobManager
@@ -79,4 +84,36 @@ public class JobExecutorFactory {
       throw new RuntimeException("Failed to create job executor: " + 
jobExecutorName, e);
     }
   }
+
+  /**
+   * Checks that the job executor class implements all the methods Gravitino 
requires. A class
+   * compiled against an older version of {@link JobExecutor} still loads, but 
calling a method it
+   * doesn't implement throws {@link AbstractMethodError} later, so it is 
rejected up front.
+   *
+   * @param jobExecutorClass The job executor class to check.
+   * @throws IllegalArgumentException If the class isn't a job executor, or 
misses a required
+   *     method.
+   */
+  @VisibleForTesting
+  static void checkJobExecutorClass(Class<?> jobExecutorClass) {
+    Preconditions.checkArgument(
+        JobExecutor.class.isAssignableFrom(jobExecutorClass),
+        "%s doesn't implement %s",
+        jobExecutorClass.getName(),
+        JobExecutor.class.getName());
+
+    Method getJobExecutionInfo;
+    try {
+      getJobExecutionInfo = jobExecutorClass.getMethod("getJobExecutionInfo", 
String.class);
+    } catch (NoSuchMethodException e) {
+      // Never happens for a JobExecutor, as the interface declares the method.
+      throw new IllegalArgumentException(e);
+    }
+    Preconditions.checkArgument(
+        !Modifier.isAbstract(getJobExecutionInfo.getModifiers()),
+        "Job executor %s doesn't implement 
JobExecutor#getJobExecutionInfo(String), which "
+            + "Gravitino uses to track the jobs. It was likely built against 
an older version of "
+            + "Gravitino, rebuild it against this version and implement the 
method.",
+        jobExecutorClass.getName());
+  }
 }
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 a2f2069d08..75b609421c 100644
--- a/core/src/main/java/org/apache/gravitino/job/JobManager.java
+++ b/core/src/main/java/org/apache/gravitino/job/JobManager.java
@@ -32,6 +32,7 @@ import java.nio.file.DirectoryStream;
 import java.nio.file.Files;
 import java.nio.file.LinkOption;
 import java.nio.file.Path;
+import java.time.Duration;
 import java.time.Instant;
 import java.util.ArrayList;
 import java.util.Arrays;
@@ -46,6 +47,7 @@ import java.util.function.Function;
 import java.util.regex.Matcher;
 import java.util.regex.Pattern;
 import java.util.stream.Collectors;
+import javax.annotation.Nullable;
 import org.apache.commons.io.FileUtils;
 import org.apache.commons.lang3.ArrayUtils;
 import org.apache.commons.lang3.StringUtils;
@@ -56,6 +58,7 @@ import org.apache.gravitino.EntityAlreadyExistsException;
 import org.apache.gravitino.EntityStore;
 import org.apache.gravitino.NameIdentifier;
 import org.apache.gravitino.Namespace;
+import org.apache.gravitino.connector.job.JobExecutionInfo;
 import org.apache.gravitino.connector.job.JobExecutor;
 import org.apache.gravitino.dto.job.JobTemplateDTO;
 import org.apache.gravitino.dto.util.DTOConverters;
@@ -86,6 +89,10 @@ public class JobManager implements JobOperationDispatcher {
 
   private static final Logger LOG = LoggerFactory.getLogger(JobManager.class);
 
+  // How far a time reported by the job executor may be ahead of Gravitino's 
clock. It tolerates the
+  // clock of the job runner being ahead, while dropping obviously invalid 
times.
+  private static final Duration MAX_REPORTED_TIME_AHEAD = Duration.ofDays(1);
+
   private static final Pattern PLACEHOLDER_PATTERN = 
Pattern.compile("\\{\\{([\\w.-]+)\\}\\}");
 
   // A job's staging directory is <stagingDir>/job-runs/job-<id>, derived from 
the job id alone:
@@ -529,6 +536,11 @@ public class JobManager implements JobOperationDispatcher {
       throw new RuntimeException("Failed to serialize the runtime job 
template", e);
     }
 
+    // The job is queued once it is submitted. Take the time before the 
submission, so that it is
+    // never later than the time the job executor reports the job started, 
which can happen right
+    // after the job is submitted.
+    Instant queuedAt = Instant.now();
+
     // Submit the job template to the job executor
     String jobExecutionId;
     try {
@@ -555,7 +567,7 @@ public class JobManager implements JobOperationDispatcher {
             .withAuditInfo(
                 AuditInfo.builder()
                     
.withCreator(PrincipalUtils.getCurrentPrincipal().getName())
-                    .withCreateTime(Instant.now())
+                    .withCreateTime(queuedAt)
                     .build())
             // A newly submitted job is queued, not started or finished yet.
             .withStartedAt(0L)
@@ -721,80 +733,25 @@ public class JobManager implements JobOperationDispatcher 
{
             // Only the job executor instance owning the job can query its 
status. The jobs
             // owned by other servers are skipped, and the jobs left behind by 
a server that has
             // exited are settled by cleanUpStagingDirs() once they expire.
-            if (jobExecutor.ownsJob(job.jobExecutionId())) {
+            if (!jobExecutor.ownsJob(job.jobExecutionId())) {
+              return;
+            }
+            // An exception escaping this scheduled task would cancel all its 
later runs, so a
+            // failure on one job must never stop the status pull of the 
others.
+            try {
               pullAndUpdateOwnedJobStatus(metalake, job);
+            } catch (RuntimeException e) {
+              LOG.error(
+                  "Failed to update the status of job {} with execution id {} 
under metalake {}",
+                  job.name(),
+                  job.jobExecutionId(),
+                  metalake,
+                  e);
             }
           });
     }
   }
 
-  private JobEntity toUpdatedStatusJobEntity(
-      JobEntity latestJobEntity, JobHandle.Status observedStatus) {
-    JobHandle.Status currentStatus = latestJobEntity.status();
-    boolean observedIsFinished = isFinishedStatus(observedStatus);
-
-    // Never regress a job out of a terminal state, and never move a 
CANCELLING job back to a
-    // non-terminal state - both would only be possible here because the 
executor status was
-    // observed against a stale snapshot of the job.
-    if (isFinishedStatus(currentStatus)
-        || (currentStatus == JobHandle.Status.CANCELLING && 
!observedIsFinished)) {
-      return latestJobEntity;
-    }
-
-    // Only a directly-observed STARTED transition is trustworthy evidence of 
when a job started.
-    // SUCCEEDED/FAILED do not prove the job ever reached STARTED: FAILED in 
particular can be
-    // reached directly from QUEUED (e.g. NoSuchJobException from the 
executor, or
-    // LocalJobExecutor failing before it records STARTED), and even for 
SUCCEEDED, backfilling
-    // startedAt from the queued time would understate queue latency and 
overstate execution
-    // duration in any derived metric. So startedAt is left unset unless a 
STARTED transition was
-    // actually observed.
-    //
-    // Only stamp startedAt on the first STARTED observation 
(latestJobEntity.startedAt() <= 0).
-    // A CANCELLING job already carries forward a real startedAt from 
cancelJob, and since
-    // cancellation is asynchronous, a poll can still observe STARTED while 
cancellation is in
-    // flight - overwriting the recorded start time with this later poll 
timestamp would lose the
-    // accurate value.
-    boolean isStarted = observedStatus == JobHandle.Status.STARTED;
-    long startedAt =
-        isStarted && latestJobEntity.startedAt() <= 0
-            ? Instant.now().toEpochMilli()
-            : latestJobEntity.startedAt();
-
-    // Preserve an already-recorded finishedAt (e.g. stamped by a concurrent 
writer) instead of
-    // overwriting it with a later poll's timestamp.
-    long finishedAt =
-        observedIsFinished
-            ? (latestJobEntity.finishedAt() > 0
-                ? latestJobEntity.finishedAt()
-                : Instant.now().toEpochMilli())
-            : latestJobEntity.finishedAt();
-
-    return JobEntity.builder()
-        .withId(latestJobEntity.id())
-        .withJobExecutionId(latestJobEntity.jobExecutionId())
-        .withJobTemplateName(latestJobEntity.jobTemplateName())
-        .withStatus(observedStatus)
-        .withNamespace(latestJobEntity.namespace())
-        .withAuditInfo(
-            AuditInfo.builder()
-                .withCreator(latestJobEntity.auditInfo().creator())
-                .withCreateTime(latestJobEntity.auditInfo().createTime())
-                
.withLastModifier(PrincipalUtils.getCurrentPrincipal().getName())
-                .withLastModifiedTime(Instant.now())
-                .build())
-        .withStartedAt(startedAt)
-        .withFinishedAt(finishedAt)
-        // The runtime job template is fixed at job creation and never changes.
-        .withRuntimeJobTemplate(latestJobEntity.runtimeJobTemplate())
-        .build();
-  }
-
-  private static boolean isFinishedStatus(JobHandle.Status status) {
-    return status == JobHandle.Status.SUCCEEDED
-        || status == JobHandle.Status.FAILED
-        || status == JobHandle.Status.CANCELLED;
-  }
-
   @VisibleForTesting
   void cleanUpStagingDirs() {
     List<String> metalakes = MetalakeManager.listInUseMetalakes(entityStore);
@@ -1211,26 +1168,28 @@ public class JobManager implements 
JobOperationDispatcher {
   }
 
   private void pullAndUpdateOwnedJobStatus(String metalake, JobEntity job) {
-    JobHandle.Status newStatus = job.status();
+    JobExecutionInfo observed;
     try {
-      newStatus = jobExecutor.getJobStatus(job.jobExecutionId());
+      observed = jobExecutor.getJobExecutionInfo(job.jobExecutionId());
       // The job was marked as CANCELLING by another server, which can't 
cancel it itself, so
       // cancel it here as the owner. This only applies to node local job 
state, other job
       // executors are cancelled directly by the server handling the request, 
and they may keep
       // reporting the job as running while cancelling it asynchronously.
       if (jobExecutor.isJobStateNodeLocal()
           && job.status() == JobHandle.Status.CANCELLING
-          && (newStatus == JobHandle.Status.QUEUED || newStatus == 
JobHandle.Status.STARTED)) {
-        newStatus = cancelOwnedJob(metalake, job);
+          && (observed.status() == JobHandle.Status.QUEUED
+              || observed.status() == JobHandle.Status.STARTED)) {
+        observed = cancelOwnedJob(metalake, job);
       }
+      observed = sanitizeExecutionInfo(job, observed);
     } catch (NoSuchJobException e) {
       // If the job is not found in the external job executor, we assume the 
job is
       // FAILED if it is not in CANCELLING status, otherwise we assume it is 
CANCELLED.
-      if (job.status() == JobHandle.Status.CANCELLING) {
-        newStatus = JobHandle.Status.CANCELLED;
-      } else {
-        newStatus = JobHandle.Status.FAILED;
-      }
+      observed =
+          JobExecutionInfo.of(
+              job.status() == JobHandle.Status.CANCELLING
+                  ? JobHandle.Status.CANCELLED
+                  : JobHandle.Status.FAILED);
       LOG.warn(
           "Job {} with execution id {} under metalake {} is not found in the "
               + "external job executor, marking it as {}. This could be due to 
the job "
@@ -1239,49 +1198,236 @@ public class JobManager implements 
JobOperationDispatcher {
           job.name(),
           job.jobExecutionId(),
           metalake,
-          newStatus);
+          observed.status());
     } catch (Exception e) {
       // Keep the job unchanged, and retry it in the next poll.
-      newStatus = job.status();
       LOG.error(
           "Failed to pull or cancel job {} by execution id {}",
           job.name(),
           job.jobExecutionId(),
           e);
+      return;
     }
 
-    if (newStatus != job.status()) {
-      // Update the job entity with new status. entityStore.update() 
re-fetches the
-      // latest entity itself right before applying the updater, so the 
transition below
-      // is derived from latestJobEntity - the state as of right before the 
write - rather
-      // than the possibly-stale `job` snapshot taken by listJobs() above. A 
concurrent
-      // writer (e.g. cancelJob(), or another poll run) may have already moved 
the job to
-      // a terminal state, into CANCELLING, or recorded a real 
startedAt/finishedAt in the
-      // gap between that snapshot and this point; the updater must not 
regress any of
-      // that using the stale snapshot's view of the world.
-      JobHandle.Status finalNewStatus = newStatus;
-      updateJobEntity(
-              metalake,
-              job,
-              latestJobEntity -> toUpdatedStatusJobEntity(latestJobEntity, 
finalNewStatus))
-          .ifPresent(
-              updated ->
-                  LOG.info(
-                      "Updated the job {} with execution id {} status to {}",
-                      job.name(),
-                      job.jobExecutionId(),
-                      updated.status()));
+    // Besides the status, the job executor may report the timestamps of the 
job later than its
+    // status, so the job is also updated when only its timestamps change. The 
job is not written
+    // at all when nothing changes.
+    if (!changesJob(job, observed)) {
+      return;
     }
+
+    // Update the job entity with new status. entityStore.update() re-fetches 
the
+    // latest entity itself right before applying the updater, so the 
transition below
+    // is derived from latestJobEntity - the state as of right before the 
write - rather
+    // than the possibly-stale `job` snapshot taken by listJobs() above. A 
concurrent
+    // writer (e.g. cancelJob(), or another poll run) may have already moved 
the job to
+    // a terminal state, into CANCELLING, or recorded a real 
startedAt/finishedAt in the
+    // gap between that snapshot and this point; the updater must not regress 
any of
+    // that using the stale snapshot's view of the world.
+    JobExecutionInfo finalObserved = observed;
+    updateJobEntity(
+            metalake,
+            job,
+            latestJobEntity -> toUpdatedStatusJobEntity(latestJobEntity, 
finalObserved))
+        .ifPresent(
+            updated ->
+                LOG.info(
+                    "Updated the job {} with execution id {} status to {}",
+                    job.name(),
+                    job.jobExecutionId(),
+                    updated.status()));
   }
 
-  private JobHandle.Status cancelOwnedJob(String metalake, JobEntity job) {
+  private JobExecutionInfo cancelOwnedJob(String metalake, JobEntity job) {
     LOG.info(
         "Cancelling job {} with execution id {} under metalake {} as it is 
marked as CANCELLING",
         job.name(),
         job.jobExecutionId(),
         metalake);
     jobExecutor.cancelJob(job.jobExecutionId());
-    return jobExecutor.getJobStatus(job.jobExecutionId());
+    return jobExecutor.getJobExecutionInfo(job.jobExecutionId());
+  }
+
+  // Drops or corrects the timestamps reported by the job executor that are 
unusable, or
+  // inconsistent
+  // with the status or with each other. A faulty job executor never fails the 
status pull. The
+  // corrections are only logged at debug level, as the same report is 
corrected on every pull.
+  private JobExecutionInfo sanitizeExecutionInfo(JobEntity job, 
JobExecutionInfo observed) {
+    Instant startedAt = usableTime(job, "started", observed.startedAt());
+    Instant finishedAt = usableTime(job, "finished", observed.finishedAt());
+
+    if (startedAt != null && observed.status() == JobHandle.Status.QUEUED) {
+      LOG.debug(
+          "Job {} with execution id {} is reported as QUEUED with a started 
time {}, ignoring "
+              + "the started time",
+          job.name(),
+          job.jobExecutionId(),
+          startedAt);
+      startedAt = null;
+    }
+    if (finishedAt != null && !isFinishedStatus(observed.status())) {
+      LOG.debug(
+          "Job {} with execution id {} is reported as {} with a finished time 
{}, ignoring the "
+              + "finished time",
+          job.name(),
+          job.jobExecutionId(),
+          observed.status(),
+          finishedAt);
+      finishedAt = null;
+    }
+
+    startedAt = notBeforeQueuedAt(job, "started", startedAt);
+    finishedAt = notBeforeQueuedAt(job, "finished", finishedAt);
+
+    if (startedAt != null && finishedAt != null && 
finishedAt.isBefore(startedAt)) {
+      // The finished time is kept, as the cleanup of finished jobs relies on 
it.
+      LOG.debug(
+          "Job {} with execution id {} is reported to finish at {} before it 
started at {}, "
+              + "ignoring the started time",
+          job.name(),
+          job.jobExecutionId(),
+          finishedAt,
+          startedAt);
+      startedAt = null;
+    }
+
+    return 
observed.toBuilder().withStartedAt(startedAt).withFinishedAt(finishedAt).build();
+  }
+
+  // A reported time is stored in epoch milliseconds, so a time that doesn't 
fit is dropped. So is a
+  // time far in the future, e.g. a sentinel like Instant.MAX, which would 
otherwise keep a finished
+  // job from ever being cleaned up.
+  @Nullable
+  private static Instant usableTime(JobEntity job, String event, @Nullable 
Instant reportedTime) {
+    if (reportedTime == null) {
+      return null;
+    }
+    boolean usable;
+    try {
+      reportedTime.toEpochMilli();
+      usable = 
!reportedTime.isAfter(Instant.now().plus(MAX_REPORTED_TIME_AHEAD));
+    } catch (ArithmeticException e) {
+      usable = false;
+    }
+    if (!usable) {
+      LOG.debug(
+          "Job {} with execution id {} is reported to be {} at an unusable 
time {}, ignoring it",
+          job.name(),
+          job.jobExecutionId(),
+          event,
+          reportedTime);
+      return null;
+    }
+    return reportedTime;
+  }
+
+  // The job executor may run on another host whose clock is behind 
Gravitino's, so a reported time
+  // can be earlier than when Gravitino queued the job. Such a time is raised 
to the queued time.
+  @Nullable
+  private Instant notBeforeQueuedAt(JobEntity job, String event, @Nullable 
Instant reportedTime) {
+    Instant queuedAt = job.auditInfo() == null ? null : 
job.auditInfo().createTime();
+    if (reportedTime == null || queuedAt == null || 
!reportedTime.isBefore(queuedAt)) {
+      return reportedTime;
+    }
+    LOG.debug(
+        "Job {} with execution id {} is reported to be {} at {}, before it was 
queued at {}, "
+            + "using the queued time instead",
+        job.name(),
+        job.jobExecutionId(),
+        event,
+        reportedTime,
+        queuedAt);
+    return queuedAt;
+  }
+
+  // Whether applying the observed execution info changes the recorded status 
or timestamps of the
+  // job. It is derived by the same logic as the update, so a value that the 
update corrects is not
+  // seen as a change on every poll.
+  private boolean changesJob(JobEntity job, JobExecutionInfo observed) {
+    JobEntity updated = toUpdatedStatusJobEntity(job, observed);
+    return updated.status() != job.status()
+        || timeInMs(updated.startedAt()) != timeInMs(job.startedAt())
+        || timeInMs(updated.finishedAt()) != timeInMs(job.finishedAt());
+  }
+
+  private JobEntity toUpdatedStatusJobEntity(JobEntity latestJobEntity, 
JobExecutionInfo observed) {
+    JobHandle.Status currentStatus = latestJobEntity.status();
+    JobHandle.Status observedStatus = observed.status();
+
+    // Never regress a job out of a terminal state, and never move a 
CANCELLING job back to a
+    // non-terminal state - both would only be possible here because the 
executor status was
+    // observed against a stale snapshot of the job.
+    if (isFinishedStatus(currentStatus)
+        || (currentStatus == JobHandle.Status.CANCELLING && 
!isFinishedStatus(observedStatus))) {
+      return latestJobEntity;
+    }
+
+    long startedAt = resolveStartedAt(latestJobEntity, observed);
+    long finishedAt = resolveFinishedAt(latestJobEntity, observed);
+    if (startedAt > 0 && finishedAt > 0 && finishedAt < startedAt) {
+      // The started time recorded by an earlier poll is later than the 
finished time reported by
+      // the job executor, so it can't be right. The finished time is kept, as 
the cleanup of
+      // finished jobs relies on it.
+      LOG.debug(
+          "Dropping the started time of job {} as it is later than its 
finished time",
+          latestJobEntity.name());
+      startedAt = 0L;
+    }
+
+    return JobEntity.builder()
+        .withId(latestJobEntity.id())
+        .withJobExecutionId(latestJobEntity.jobExecutionId())
+        .withJobTemplateName(latestJobEntity.jobTemplateName())
+        .withStatus(observedStatus)
+        .withNamespace(latestJobEntity.namespace())
+        .withAuditInfo(
+            AuditInfo.builder()
+                .withCreator(latestJobEntity.auditInfo().creator())
+                .withCreateTime(latestJobEntity.auditInfo().createTime())
+                
.withLastModifier(PrincipalUtils.getCurrentPrincipal().getName())
+                .withLastModifiedTime(Instant.now())
+                .build())
+        .withStartedAt(startedAt)
+        .withFinishedAt(finishedAt)
+        // The runtime job template is fixed at job creation and never changes.
+        .withRuntimeJobTemplate(latestJobEntity.runtimeJobTemplate())
+        .build();
+  }
+
+  private static long resolveStartedAt(JobEntity latestJobEntity, 
JobExecutionInfo observed) {
+    // The time reported by the job executor is when the job actually started, 
so it also replaces
+    // a time recorded by an earlier poll.
+    if (observed.startedAt() != null) {
+      return observed.startedAt().toEpochMilli();
+    }
+
+    // The job executor doesn't report when the job started. Only a 
directly-observed STARTED
+    // transition is then trustworthy evidence of when a job started. 
SUCCEEDED/FAILED do not prove
+    // the job ever reached STARTED: FAILED in particular can be reached 
directly from QUEUED (e.g.
+    // NoSuchJobException from the executor), and even for SUCCEEDED, 
backfilling startedAt from
+    // the queued time would understate queue latency and overstate execution 
duration in any
+    // derived metric. So startedAt is left unset unless a STARTED transition 
was actually
+    // observed, and it is only stamped on the first STARTED observation.
+    long recordedStartedAt = timeInMs(latestJobEntity.startedAt());
+    if (observed.status() == JobHandle.Status.STARTED && recordedStartedAt <= 
0) {
+      return Instant.now().toEpochMilli();
+    }
+    return recordedStartedAt;
+  }
+
+  private static long resolveFinishedAt(JobEntity latestJobEntity, 
JobExecutionInfo observed) {
+    long recordedFinishedAt = timeInMs(latestJobEntity.finishedAt());
+    if (!isFinishedStatus(observed.status())) {
+      return recordedFinishedAt;
+    }
+    if (observed.finishedAt() != null) {
+      return observed.finishedAt().toEpochMilli();
+    }
+
+    // The job executor doesn't report when the job finished. Preserve an 
already-recorded
+    // finishedAt (e.g. stamped by a concurrent writer) instead of overwriting 
it with a later
+    // poll's timestamp.
+    return recordedFinishedAt > 0 ? recordedFinishedAt : 
Instant.now().toEpochMilli();
   }
 
   private Optional<JobEntity> updateJobEntity(
@@ -1340,9 +1486,10 @@ public class JobManager implements 
JobOperationDispatcher {
               // another retention time like any other finished job before 
being cleaned up.
               return toUpdatedStatusJobEntity(
                   latestJobEntity,
-                  latestJobEntity.status() == JobHandle.Status.CANCELLING
-                      ? JobHandle.Status.CANCELLED
-                      : JobHandle.Status.FAILED);
+                  JobExecutionInfo.of(
+                      latestJobEntity.status() == JobHandle.Status.CANCELLING
+                          ? JobHandle.Status.CANCELLED
+                          : JobHandle.Status.FAILED));
             })
         .filter(expiredJob -> expired.get())
         .ifPresent(
@@ -1371,4 +1518,14 @@ public class JobManager implements 
JobOperationDispatcher {
             : auditInfo.createTime();
     return lastUpdatedTime == null ? 0L : lastUpdatedTime.toEpochMilli();
   }
+
+  private static boolean isFinishedStatus(JobHandle.Status status) {
+    return status == JobHandle.Status.SUCCEEDED
+        || status == JobHandle.Status.FAILED
+        || status == JobHandle.Status.CANCELLED;
+  }
+
+  private static long timeInMs(@Nullable Long time) {
+    return time == null ? 0L : time;
+  }
 }
diff --git 
a/core/src/main/java/org/apache/gravitino/job/local/LocalJobExecutor.java 
b/core/src/main/java/org/apache/gravitino/job/local/LocalJobExecutor.java
index 296a6d89de..a8f83702c8 100644
--- a/core/src/main/java/org/apache/gravitino/job/local/LocalJobExecutor.java
+++ b/core/src/main/java/org/apache/gravitino/job/local/LocalJobExecutor.java
@@ -49,6 +49,7 @@ import java.nio.file.NoSuchFileException;
 import java.nio.file.Path;
 import java.nio.file.Paths;
 import java.nio.file.StandardOpenOption;
+import java.time.Instant;
 import java.util.Arrays;
 import java.util.List;
 import java.util.Map;
@@ -66,6 +67,7 @@ import java.util.regex.Pattern;
 import javax.annotation.Nullable;
 import org.apache.commons.lang3.StringUtils;
 import org.apache.commons.lang3.tuple.Pair;
+import org.apache.gravitino.connector.job.JobExecutionInfo;
 import org.apache.gravitino.connector.job.JobExecutor;
 import org.apache.gravitino.exceptions.NoSuchJobException;
 import org.apache.gravitino.job.JobHandle;
@@ -93,8 +95,6 @@ public class LocalJobExecutor implements JobExecutor {
 
   private static final String LOCAL_JOB_PREFIX = "local-job-";
 
-  private static final long UNEXPIRED_TIME_IN_MS = -1L;
-
   // Contains '.' and '-', which MetalakeNormalizeDispatcher rejects in 
metalake names on create and
   // rename, so it doesn't collide with a metalake's staging directory. A 
metalake created before
   // that check could still collide: the cleanup then skips its directories 
rather than failing.
@@ -137,10 +137,10 @@ public class LocalJobExecutor implements JobExecutor {
 
   private ExecutorService jobPollingExecutorService;
 
-  // The job status map to keep track of the status of each job. In the 
meantime, the job status
-  // will be stored in the entity store, so we will clean the finished, 
cancelled and failed jobs
-  // from the map periodically to save the memory.
-  private Map<String, Pair<JobHandle.Status, Long>> jobStatus;
+  // The map to keep track of the status and the timestamps of each job. In 
the meantime, the job
+  // state will be stored in the entity store, so we will clean the finished, 
cancelled and failed
+  // jobs from the map periodically to save the memory.
+  private Map<String, JobExecutionInfo> jobInfos;
   private final Object lock = new Object();
 
   private long jobStatusKeepTimeInMs;
@@ -244,7 +244,7 @@ public class LocalJobExecutor implements JobExecutor {
     threadPoolExecutor.allowCoreThreadTimeOut(true);
     this.jobExecutorService = threadPoolExecutor;
 
-    this.jobStatus = Maps.newHashMap();
+    this.jobInfos = Maps.newHashMap();
 
     this.jobStatusKeepTimeInMs =
         configs.containsKey(JOB_STATUS_KEEP_TIME_MS)
@@ -314,7 +314,7 @@ public class LocalJobExecutor implements JobExecutor {
         throw new IllegalStateException("Waiting queue is full, cannot submit 
job: " + jobTemplate);
       }
 
-      jobStatus.put(newJobId, Pair.of(JobHandle.Status.QUEUED, 
UNEXPIRED_TIME_IN_MS));
+      jobInfos.put(newJobId, JobExecutionInfo.of(JobHandle.Status.QUEUED));
     }
 
     // Written only once the job is accepted, so a rejected submission leaves 
no index behind.
@@ -323,55 +323,52 @@ public class LocalJobExecutor implements JobExecutor {
   }
 
   @Override
-  public JobHandle.Status getJobStatus(String jobId) throws NoSuchJobException 
{
+  public JobExecutionInfo getJobExecutionInfo(String jobId) throws 
NoSuchJobException {
     synchronized (lock) {
-      if (!jobStatus.containsKey(jobId)) {
+      JobExecutionInfo info = jobInfos.get(jobId);
+      if (info == null) {
         throw new NoSuchJobException("No job found with ID: %s", jobId);
       }
-      LOG.debug(
-          "Get status {} and finished time {} for job {}",
-          jobStatus.get(jobId).getLeft(),
-          jobStatus.get(jobId).getRight(),
-          jobId);
-      return jobStatus.get(jobId).getLeft();
+      LOG.debug("Get {} for job {}", info, jobId);
+      return info;
     }
   }
 
   @Override
   public void cancelJob(String jobId) throws NoSuchJobException {
     synchronized (lock) {
-      if (!jobStatus.containsKey(jobId)) {
+      JobExecutionInfo info = jobInfos.get(jobId);
+      if (info == null) {
         throw new NoSuchJobException("No job found with ID: %s", jobId);
       }
 
-      Pair<JobHandle.Status, Long> statusPair = jobStatus.get(jobId);
-      if (statusPair.getLeft() == JobHandle.Status.SUCCEEDED
-          || statusPair.getLeft() == JobHandle.Status.FAILED
-          || statusPair.getLeft() == JobHandle.Status.CANCELLED) {
+      if (info.status() == JobHandle.Status.SUCCEEDED
+          || info.status() == JobHandle.Status.FAILED
+          || info.status() == JobHandle.Status.CANCELLED) {
         LOG.warn("Job {} is already completed or cancelled, no action taken", 
jobId);
         return;
       }
 
-      if (statusPair.getLeft() == JobHandle.Status.CANCELLING) {
+      if (info.status() == JobHandle.Status.CANCELLING) {
         LOG.warn("Job {} is already being cancelled, no action taken", jobId);
         return;
       }
 
       // If the job is queued.
-      if (statusPair.getLeft() == JobHandle.Status.QUEUED) {
+      if (info.status() == JobHandle.Status.QUEUED) {
         waitingQueue.removeIf(p -> p.getLeft().equals(jobId));
-        jobStatus.put(jobId, Pair.of(JobHandle.Status.CANCELLED, 
System.currentTimeMillis()));
+        finishJob(jobId, JobHandle.Status.CANCELLED);
         LOG.info("Job {} is cancelled from the waiting queue", jobId);
         return;
       }
 
-      if (statusPair.getLeft() == JobHandle.Status.STARTED) {
+      if (info.status() == JobHandle.Status.STARTED) {
         Process process = runningProcesses.get(jobId);
         if (process != null) {
           process.destroy();
         }
         LOG.info("Job {} is cancelling while running", jobId);
-        jobStatus.put(jobId, Pair.of(JobHandle.Status.CANCELLING, 
UNEXPIRED_TIME_IN_MS));
+        jobInfos.put(jobId, 
info.toBuilder().withStatus(JobHandle.Status.CANCELLING).build());
       }
     }
   }
@@ -423,7 +420,7 @@ public class LocalJobExecutor implements JobExecutor {
     // Stop the job status and output index cleanup executors
     jobStatusCleanupExecutor.shutdownNow();
     outputIndexCleanupExecutor.shutdownNow();
-    jobStatus.clear();
+    jobInfos.clear();
   }
 
   public void runJob(Pair<String, JobTemplate> jobPair) {
@@ -434,8 +431,8 @@ public class LocalJobExecutor implements JobExecutor {
       Process process;
       synchronized (lock) {
         // This happens when the job is cancelled before it starts.
-        Pair<JobHandle.Status, Long> statusPair = jobStatus.get(jobId);
-        if (statusPair == null || statusPair.getLeft() != 
JobHandle.Status.QUEUED) {
+        JobExecutionInfo info = jobInfos.get(jobId);
+        if (info == null || info.status() != JobHandle.Status.QUEUED) {
           LOG.warn("Job {} is not in QUEUED state, cannot start it", jobId);
           return;
         }
@@ -443,7 +440,7 @@ public class LocalJobExecutor implements JobExecutor {
         LocalProcessBuilder processBuilder = 
LocalProcessBuilder.create(jobTemplate, configs);
         process = processBuilder.start();
         runningProcesses.put(jobId, process);
-        jobStatus.put(jobId, Pair.of(JobHandle.Status.STARTED, 
UNEXPIRED_TIME_IN_MS));
+        jobInfos.put(jobId, info.started(Instant.now()));
       }
 
       LOG.info("Starting job: {}", jobId);
@@ -452,17 +449,17 @@ public class LocalJobExecutor implements JobExecutor {
       if (exitCode == 0) {
         LOG.info("Job {} completed successfully", jobId);
         synchronized (lock) {
-          jobStatus.put(jobId, Pair.of(JobHandle.Status.SUCCEEDED, 
System.currentTimeMillis()));
+          finishJob(jobId, JobHandle.Status.SUCCEEDED);
         }
       } else {
         synchronized (lock) {
-          JobHandle.Status oldStatus = jobStatus.get(jobId).getLeft();
+          JobHandle.Status oldStatus = jobInfos.get(jobId).status();
           if (oldStatus == JobHandle.Status.CANCELLING) {
             LOG.info("Job {} was cancelled while running with exit code: {}", 
jobId, exitCode);
-            jobStatus.put(jobId, Pair.of(JobHandle.Status.CANCELLED, 
System.currentTimeMillis()));
+            finishJob(jobId, JobHandle.Status.CANCELLED);
           } else if (oldStatus == JobHandle.Status.STARTED) {
             LOG.warn("Job {} failed after starting with exit code: {}", jobId, 
exitCode);
-            jobStatus.put(jobId, Pair.of(JobHandle.Status.FAILED, 
System.currentTimeMillis()));
+            finishJob(jobId, JobHandle.Status.FAILED);
           }
         }
       }
@@ -471,8 +468,7 @@ public class LocalJobExecutor implements JobExecutor {
       LOG.error("Error while executing job", e);
       // If an error occurs, we should mark the job as failed
       synchronized (lock) {
-        String jobId = jobPair.getLeft();
-        jobStatus.put(jobId, Pair.of(JobHandle.Status.FAILED, 
System.currentTimeMillis()));
+        finishJob(jobPair.getLeft(), JobHandle.Status.FAILED);
       }
     }
 
@@ -507,12 +503,13 @@ public class LocalJobExecutor implements JobExecutor {
     long currentTime = System.currentTimeMillis();
 
     synchronized (lock) {
-      jobStatus
-          .entrySet()
+      // Only the finished jobs have a finished time, the others are never 
removed.
+      jobInfos
+          .values()
           .removeIf(
-              entry ->
-                  entry.getValue().getRight() != UNEXPIRED_TIME_IN_MS
-                      && (currentTime - entry.getValue().getRight()) >= 
jobStatusKeepTimeInMs);
+              info ->
+                  info.finishedAt() != null
+                      && (currentTime - info.finishedAt().toEpochMilli()) >= 
jobStatusKeepTimeInMs);
     }
   }
 
@@ -561,6 +558,14 @@ public class LocalJobExecutor implements JobExecutor {
     }
   }
 
+  // Must be called with the lock held. The started time, if any, is carried 
forward.
+  private void finishJob(String jobId, JobHandle.Status status) {
+    // The job may be missing if the executor is closed concurrently.
+    JobExecutionInfo info =
+        jobInfos.getOrDefault(jobId, 
JobExecutionInfo.of(JobHandle.Status.QUEUED));
+    jobInfos.put(jobId, info.finished(status, Instant.now()));
+  }
+
   private List<String> getJobOutput(String jobId, String fileName, int 
maxLines, int maxBytes) {
     Path workingDir = locateWorkingDir(jobId);
     if (workingDir == null || !isRealWorkingDir(jobId, workingDir)) {
diff --git 
a/core/src/main/java/org/apache/gravitino/listener/api/info/JobInfo.java 
b/core/src/main/java/org/apache/gravitino/listener/api/info/JobInfo.java
index cf6b9be02f..d5e68b1f08 100644
--- a/core/src/main/java/org/apache/gravitino/listener/api/info/JobInfo.java
+++ b/core/src/main/java/org/apache/gravitino/listener/api/info/JobInfo.java
@@ -160,7 +160,11 @@ public final class JobInfo {
   /**
    * Returns the time when the job started execution.
    *
-   * @return the started time of the job, or null if the job has not started 
execution yet
+   * <p>A finished job may also have no started time, if the job executor 
doesn't report when the
+   * job started and Gravitino didn't observe the job running.
+   *
+   * @return the started time of the job, or null if the job has not started 
execution yet, or the
+   *     started time is unknown
    */
   @Nullable
   public Instant startedAt() {
diff --git 
a/core/src/test/java/org/apache/gravitino/connector/job/TestJobExecutionInfo.java
 
b/core/src/test/java/org/apache/gravitino/connector/job/TestJobExecutionInfo.java
new file mode 100644
index 0000000000..6692090b33
--- /dev/null
+++ 
b/core/src/test/java/org/apache/gravitino/connector/job/TestJobExecutionInfo.java
@@ -0,0 +1,139 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *  http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.gravitino.connector.job;
+
+import java.io.IOException;
+import java.time.Instant;
+import java.util.Map;
+import org.apache.gravitino.job.JobHandle;
+import org.apache.gravitino.job.JobTemplate;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+public class TestJobExecutionInfo {
+
+  @Test
+  public void testOf() {
+    JobExecutionInfo info = JobExecutionInfo.of(JobHandle.Status.QUEUED);
+    Assertions.assertEquals(JobHandle.Status.QUEUED, info.status());
+    Assertions.assertNull(info.startedAt());
+    Assertions.assertNull(info.finishedAt());
+  }
+
+  @Test
+  public void testBuilder() {
+    Instant startedAt = Instant.ofEpochMilli(1000L);
+    Instant finishedAt = Instant.ofEpochMilli(2000L);
+    JobExecutionInfo info =
+        JobExecutionInfo.builder()
+            .withStatus(JobHandle.Status.SUCCEEDED)
+            .withStartedAt(startedAt)
+            .withFinishedAt(finishedAt)
+            .build();
+
+    Assertions.assertEquals(JobHandle.Status.SUCCEEDED, info.status());
+    Assertions.assertEquals(startedAt, info.startedAt());
+    Assertions.assertEquals(finishedAt, info.finishedAt());
+    Assertions.assertEquals(
+        info,
+        JobExecutionInfo.builder()
+            .withStatus(JobHandle.Status.SUCCEEDED)
+            .withStartedAt(startedAt)
+            .withFinishedAt(finishedAt)
+            .build());
+    Assertions.assertNotEquals(info, 
JobExecutionInfo.of(JobHandle.Status.SUCCEEDED));
+
+    Assertions.assertThrows(NullPointerException.class, () -> 
JobExecutionInfo.of(null));
+    Assertions.assertThrows(
+        NullPointerException.class,
+        () -> JobExecutionInfo.builder().withStartedAt(startedAt).build());
+    Assertions.assertEquals(
+        JobExecutionInfo.builder()
+            .withStatus(JobHandle.Status.FAILED)
+            .withStartedAt(startedAt)
+            .withFinishedAt(finishedAt)
+            .build(),
+        info.toBuilder().withStatus(JobHandle.Status.FAILED).build());
+  }
+
+  @Test
+  public void testStartedAndFinished() {
+    Instant startedAt = Instant.ofEpochMilli(1000L);
+    Instant finishedAt = Instant.ofEpochMilli(2000L);
+
+    JobExecutionInfo started = 
JobExecutionInfo.of(JobHandle.Status.QUEUED).started(startedAt);
+    Assertions.assertEquals(JobHandle.Status.STARTED, started.status());
+    Assertions.assertEquals(startedAt, started.startedAt());
+    Assertions.assertNull(started.finishedAt());
+
+    // The started time is carried forward to the finished job.
+    JobExecutionInfo succeeded = started.finished(JobHandle.Status.SUCCEEDED, 
finishedAt);
+    Assertions.assertEquals(JobHandle.Status.SUCCEEDED, succeeded.status());
+    Assertions.assertEquals(startedAt, succeeded.startedAt());
+    Assertions.assertEquals(finishedAt, succeeded.finishedAt());
+
+    // A job cancelled before it started has no started time.
+    JobExecutionInfo cancelled =
+        JobExecutionInfo.of(JobHandle.Status.QUEUED)
+            .finished(JobHandle.Status.CANCELLED, finishedAt);
+    Assertions.assertEquals(JobHandle.Status.CANCELLED, cancelled.status());
+    Assertions.assertNull(cancelled.startedAt());
+    Assertions.assertEquals(finishedAt, cancelled.finishedAt());
+
+    Assertions.assertThrows(IllegalArgumentException.class, () -> 
started.started(null));
+    Assertions.assertThrows(
+        IllegalArgumentException.class,
+        () -> started.finished(JobHandle.Status.CANCELLING, finishedAt));
+    Assertions.assertThrows(
+        IllegalArgumentException.class, () -> 
started.finished(JobHandle.Status.FAILED, null));
+  }
+
+  @Test
+  public void testDefaultGetJobStatus() {
+    // getJobStatus is derived from the execution info reported by the job 
executor.
+    JobExecutionInfo info =
+        JobExecutionInfo.of(JobHandle.Status.QUEUED)
+            .started(Instant.ofEpochMilli(1000L))
+            .finished(JobHandle.Status.SUCCEEDED, Instant.ofEpochMilli(2000L));
+    JobExecutor jobExecutor =
+        new JobExecutor() {
+          @Override
+          public void initialize(Map<String, String> configs) {}
+
+          @Override
+          public String submitJob(JobTemplate jobTemplate) {
+            return "job-1";
+          }
+
+          @Override
+          public JobExecutionInfo getJobExecutionInfo(String jobId) {
+            return info;
+          }
+
+          @Override
+          public void cancelJob(String jobId) {}
+
+          @Override
+          public void close() throws IOException {}
+        };
+
+    Assertions.assertEquals(JobHandle.Status.SUCCEEDED, 
jobExecutor.getJobStatus("job-1"));
+  }
+}
diff --git 
a/core/src/test/java/org/apache/gravitino/job/TestJobExecutorFactory.java 
b/core/src/test/java/org/apache/gravitino/job/TestJobExecutorFactory.java
index b8f1a8399c..d767723c8b 100644
--- a/core/src/test/java/org/apache/gravitino/job/TestJobExecutorFactory.java
+++ b/core/src/test/java/org/apache/gravitino/job/TestJobExecutorFactory.java
@@ -21,11 +21,17 @@ package org.apache.gravitino.job;
 import com.google.common.collect.ImmutableMap;
 import java.io.File;
 import java.io.IOException;
+import java.net.URL;
+import java.net.URLClassLoader;
+import java.nio.charset.StandardCharsets;
 import java.nio.file.Files;
 import java.util.Map;
+import javax.tools.JavaCompiler;
+import javax.tools.ToolProvider;
 import org.apache.commons.io.FileUtils;
 import org.apache.gravitino.Config;
 import org.apache.gravitino.Configs;
+import org.apache.gravitino.connector.job.JobExecutionInfo;
 import org.apache.gravitino.connector.job.JobExecutor;
 import org.apache.gravitino.exceptions.NoSuchJobException;
 import org.apache.gravitino.job.local.LocalJobExecutor;
@@ -128,6 +134,106 @@ public class TestJobExecutorFactory {
     }
   }
 
+  @Test
+  public void testRejectJobExecutorBuiltAgainstOldSpi() throws Exception {
+    // A job executor plugin built before getJobExecutionInfo was added to the 
SPI only
+    // implements getJobStatus. Loaded against the current SPI, it would throw 
AbstractMethodError
+    // on every status pull, so it must be rejected when the job executor is 
created.
+    Class<?> oldJobExecutorClass = compileAgainstOldSpi();
+    try {
+      JobExecutor oldJobExecutor =
+          (JobExecutor) 
oldJobExecutorClass.getDeclaredConstructor().newInstance();
+      Assertions.assertThrows(
+          AbstractMethodError.class, () -> 
oldJobExecutor.getJobExecutionInfo("job-1"));
+
+      IllegalArgumentException e =
+          Assertions.assertThrows(
+              IllegalArgumentException.class,
+              () -> 
JobExecutorFactory.checkJobExecutorClass(oldJobExecutorClass));
+      Assertions.assertTrue(e.getMessage().contains("getJobExecutionInfo"), 
e.getMessage());
+    } finally {
+      ((URLClassLoader) oldJobExecutorClass.getClassLoader()).close();
+    }
+
+    Assertions.assertDoesNotThrow(
+        () -> 
JobExecutorFactory.checkJobExecutorClass(RecordingJobExecutor.class));
+    Assertions.assertThrows(
+        IllegalArgumentException.class,
+        () -> JobExecutorFactory.checkJobExecutorClass(String.class));
+  }
+
+  // Compiles a job executor against the SPI as it was before 
getJobExecutionInfo was added, and
+  // loads it against the current SPI, like a plugin jar built for an older 
Gravitino version.
+  private Class<?> compileAgainstOldSpi() throws IOException {
+    File sourceDir = new File(testDir, "old-spi-src");
+    File classDir = new File(testDir, "old-spi-classes");
+    Assertions.assertTrue(classDir.mkdirs());
+
+    File oldSpi = new File(sourceDir, 
"org/apache/gravitino/connector/job/JobExecutor.java");
+    FileUtils.writeStringToFile(
+        oldSpi,
+        String.join(
+            "\n",
+            "package org.apache.gravitino.connector.job;",
+            "import java.util.Map;",
+            "import org.apache.gravitino.job.JobHandle;",
+            "import org.apache.gravitino.job.JobTemplate;",
+            "public interface JobExecutor extends java.io.Closeable {",
+            "  void initialize(Map<String, String> configs);",
+            "  String submitJob(JobTemplate jobTemplate);",
+            "  JobHandle.Status getJobStatus(String jobId);",
+            "  void cancelJob(String jobId);",
+            "}"),
+        StandardCharsets.UTF_8);
+    File oldExecutor = new File(sourceDir, "com/example/OldJobExecutor.java");
+    FileUtils.writeStringToFile(
+        oldExecutor,
+        String.join(
+            "\n",
+            "package com.example;",
+            "import java.util.Map;",
+            "import org.apache.gravitino.connector.job.JobExecutor;",
+            "import org.apache.gravitino.job.JobHandle;",
+            "import org.apache.gravitino.job.JobTemplate;",
+            "public class OldJobExecutor implements JobExecutor {",
+            "  public void initialize(Map<String, String> configs) {}",
+            "  public String submitJob(JobTemplate jobTemplate) { return 
\"job-1\"; }",
+            "  public JobHandle.Status getJobStatus(String jobId) {",
+            "    return JobHandle.Status.SUCCEEDED;",
+            "  }",
+            "  public void cancelJob(String jobId) {}",
+            "  public void close() {}",
+            "}"),
+        StandardCharsets.UTF_8);
+
+    JavaCompiler compiler = ToolProvider.getSystemJavaCompiler();
+    Assertions.assertNotNull(compiler, "The tests must run on a JDK");
+    int result =
+        compiler.run(
+            null,
+            null,
+            null,
+            "-classpath",
+            System.getProperty("java.class.path"),
+            "-d",
+            classDir.getAbsolutePath(),
+            oldSpi.getAbsolutePath(),
+            oldExecutor.getAbsolutePath());
+    Assertions.assertEquals(0, result);
+
+    // Only the plugin class is loaded from the compiled classes: the class 
loader delegates to its
+    // parent first, so JobExecutor resolves to the current SPI.
+    URLClassLoader classLoader =
+        new URLClassLoader(
+            new URL[] {classDir.toURI().toURL()}, 
TestJobExecutorFactory.class.getClassLoader());
+    try {
+      return classLoader.loadClass("com.example.OldJobExecutor");
+    } catch (ClassNotFoundException e) {
+      classLoader.close();
+      throw new IOException(e);
+    }
+  }
+
   /** A user's subclass of the local job executor, inheriting its 
initialization. */
   public static class CustomLocalJobExecutor extends LocalJobExecutor {}
 
@@ -147,7 +253,7 @@ public class TestJobExecutorFactory {
     }
 
     @Override
-    public JobHandle.Status getJobStatus(String jobId) throws 
NoSuchJobException {
+    public JobExecutionInfo getJobExecutionInfo(String jobId) throws 
NoSuchJobException {
       throw new NoSuchJobException("No job found with ID: %s", jobId);
     }
 
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 d19c274c02..a1b0f9f485 100644
--- a/core/src/test/java/org/apache/gravitino/job/TestJobManager.java
+++ b/core/src/test/java/org/apache/gravitino/job/TestJobManager.java
@@ -43,6 +43,7 @@ import java.net.InetSocketAddress;
 import java.nio.charset.StandardCharsets;
 import java.nio.file.Files;
 import java.nio.file.Path;
+import java.time.Duration;
 import java.time.Instant;
 import java.time.temporal.ChronoUnit;
 import java.util.Collections;
@@ -53,6 +54,7 @@ import java.util.Random;
 import java.util.UUID;
 import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
 import java.util.function.Function;
 import java.util.stream.Collectors;
 import javax.annotation.Nullable;
@@ -67,6 +69,7 @@ import org.apache.gravitino.EntityStore;
 import org.apache.gravitino.GravitinoEnv;
 import org.apache.gravitino.NameIdentifier;
 import org.apache.gravitino.Namespace;
+import org.apache.gravitino.connector.job.JobExecutionInfo;
 import org.apache.gravitino.connector.job.JobExecutor;
 import org.apache.gravitino.dto.job.JobTemplateDTO;
 import org.apache.gravitino.dto.job.ShellJobTemplateDTO;
@@ -792,6 +795,34 @@ public class TestJobManager {
         () -> jobManager.runJob(metalake, "shell_job", 
Collections.emptyMap()));
   }
 
+  @Test
+  public void testRunJobQueuesJobBeforeSubmission() throws IOException {
+    mockedMetalake
+        .when(() -> MetalakeManager.checkMetalake(metalakeIdent, entityStore))
+        .thenAnswer(a -> null);
+    JobTemplateEntity shellJobTemplate =
+        newShellJobTemplateEntity("shell_job", "A shell job template");
+    when(jobManager.getJobTemplate(metalake, 
shellJobTemplate.name())).thenReturn(shellJobTemplate);
+    doNothing().when(entityStore).put(any(JobEntity.class), anyBoolean());
+
+    // The job executor may start the job right after it is submitted, so the 
queued time must be
+    // taken before the submission to never be later than the reported started 
time.
+    AtomicReference<Instant> submittedAt = new AtomicReference<>();
+    when(jobExecutor.submitJob(any()))
+        .thenAnswer(
+            invocation -> {
+              Thread.sleep(5);
+              submittedAt.set(Instant.now());
+              return "job_execution_id_for_test";
+            });
+
+    JobEntity jobEntity = jobManager.runJob(metalake, "shell_job", 
Collections.emptyMap());
+
+    Assertions.assertTrue(
+        jobEntity.auditInfo().createTime().isBefore(submittedAt.get()),
+        "The queued time should be taken before the job is submitted");
+  }
+
   @Test
   public void testRunJobPropagatesJobExecutorRejection() throws IOException {
     mockedMetalake
@@ -1113,13 +1144,15 @@ public class TestJobManager {
 
     when(jobManager.listJobs(metalake, 
Optional.empty())).thenReturn(ImmutableList.of(job));
 
-    
when(jobExecutor.getJobStatus(job.jobExecutionId())).thenReturn(JobHandle.Status.QUEUED);
+    when(jobExecutor.getJobExecutionInfo(job.jobExecutionId()))
+        .thenReturn(JobExecutionInfo.of(JobHandle.Status.QUEUED));
     Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
     verify(entityStore, never())
         .update(any(), eq(JobEntity.class), eq(Entity.EntityType.JOB), any());
 
     stubEntityStoreUpdateToApply(job);
-    
when(jobExecutor.getJobStatus(job.jobExecutionId())).thenReturn(JobHandle.Status.SUCCEEDED);
+    when(jobExecutor.getJobExecutionInfo(job.jobExecutionId()))
+        .thenReturn(JobExecutionInfo.of(JobHandle.Status.SUCCEEDED));
     Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
 
     // Once a job transitions to a terminal status, finishedAt must be set.
@@ -1159,7 +1192,8 @@ public class TestJobManager {
         .when(() -> MetalakeManager.listInUseMetalakes(entityStore))
         .thenReturn(ImmutableList.of(metalake));
     when(jobManager.listJobs(metalake, 
Optional.empty())).thenReturn(ImmutableList.of(job));
-    
when(jobExecutor.getJobStatus(job.jobExecutionId())).thenReturn(JobHandle.Status.SUCCEEDED);
+    when(jobExecutor.getJobExecutionInfo(job.jobExecutionId()))
+        .thenReturn(JobExecutionInfo.of(JobHandle.Status.SUCCEEDED));
     stubEntityStoreUpdateToApply(job);
 
     Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
@@ -1202,7 +1236,8 @@ public class TestJobManager {
 
     // QUEUED -> STARTED: startedAt must be set, finishedAt must remain unset.
     stubEntityStoreUpdateToApply(job);
-    
when(jobExecutor.getJobStatus(job.jobExecutionId())).thenReturn(JobHandle.Status.STARTED);
+    when(jobExecutor.getJobExecutionInfo(job.jobExecutionId()))
+        .thenReturn(JobExecutionInfo.of(JobHandle.Status.STARTED));
     Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
 
     JobEntity startedJob = captureUpdatedJobEntity(job);
@@ -1216,7 +1251,8 @@ public class TestJobManager {
     Mockito.clearInvocations(entityStore);
     stubEntityStoreUpdateToApply(startedJob);
     when(jobManager.listJobs(metalake, 
Optional.empty())).thenReturn(ImmutableList.of(startedJob));
-    
when(jobExecutor.getJobStatus(job.jobExecutionId())).thenReturn(JobHandle.Status.SUCCEEDED);
+    when(jobExecutor.getJobExecutionInfo(job.jobExecutionId()))
+        .thenReturn(JobExecutionInfo.of(JobHandle.Status.SUCCEEDED));
     Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
 
     JobEntity finishedJob = captureUpdatedJobEntity(startedJob);
@@ -1262,8 +1298,8 @@ public class TestJobManager {
 
     when(jobManager.listJobs(metalake, 
Optional.empty())).thenReturn(ImmutableList.of(queuedJob));
     stubEntityStoreUpdateToApply(queuedJob);
-    when(jobExecutor.getJobStatus(queuedJob.jobExecutionId()))
-        .thenReturn(JobHandle.Status.SUCCEEDED);
+    when(jobExecutor.getJobExecutionInfo(queuedJob.jobExecutionId()))
+        .thenReturn(JobExecutionInfo.of(JobHandle.Status.SUCCEEDED));
     Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
 
     JobEntity succeededJob = captureUpdatedJobEntity(queuedJob);
@@ -1306,8 +1342,8 @@ public class TestJobManager {
     when(jobManager.listJobs(metalake, Optional.empty()))
         .thenReturn(ImmutableList.of(cancellingJob));
     stubEntityStoreUpdateToApply(cancellingJob);
-    when(jobExecutor.getJobStatus(cancellingJob.jobExecutionId()))
-        .thenReturn(JobHandle.Status.CANCELLED);
+    when(jobExecutor.getJobExecutionInfo(cancellingJob.jobExecutionId()))
+        .thenReturn(JobExecutionInfo.of(JobHandle.Status.CANCELLED));
     Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
 
     JobEntity cancelledJob = captureUpdatedJobEntity(cancellingJob);
@@ -1365,8 +1401,8 @@ public class TestJobManager {
     stubEntityStoreUpdateToApply(latestSucceeded);
     // The stale QUEUED snapshot leads the poll to observe (and try to apply) 
FAILED - a
     // different terminal status than the one the job has actually already 
settled into.
-    when(jobExecutor.getJobStatus(queuedSnapshot.jobExecutionId()))
-        .thenReturn(JobHandle.Status.FAILED);
+    when(jobExecutor.getJobExecutionInfo(queuedSnapshot.jobExecutionId()))
+        .thenReturn(JobExecutionInfo.of(JobHandle.Status.FAILED));
     Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
 
     JobEntity result = captureUpdatedJobEntity(latestSucceeded);
@@ -1410,8 +1446,8 @@ public class TestJobManager {
     when(jobManager.listJobs(metalake, Optional.empty()))
         .thenReturn(ImmutableList.of(queuedSnapshot));
     stubEntityStoreUpdateToApply(latestCancelling);
-    when(jobExecutor.getJobStatus(queuedSnapshot.jobExecutionId()))
-        .thenReturn(JobHandle.Status.STARTED);
+    when(jobExecutor.getJobExecutionInfo(queuedSnapshot.jobExecutionId()))
+        .thenReturn(JobExecutionInfo.of(JobHandle.Status.STARTED));
     Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
 
     JobEntity result = captureUpdatedJobEntity(latestCancelling);
@@ -1436,7 +1472,7 @@ public class TestJobManager {
 
     Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
 
-    verify(jobExecutor, never()).getJobStatus(any());
+    verify(jobExecutor, never()).getJobExecutionInfo(any());
     verify(entityStore, never())
         .update(any(), eq(JobEntity.class), eq(Entity.EntityType.JOB), any());
   }
@@ -1449,8 +1485,10 @@ public class TestJobManager {
         newJobEntity("local-job-mine-1", JobHandle.Status.CANCELLING, 
Instant.now(), null);
     mockListActiveJobs(startedJob);
     when(jobExecutor.isJobStateNodeLocal()).thenReturn(true);
-    when(jobExecutor.getJobStatus(startedJob.jobExecutionId()))
-        .thenReturn(JobHandle.Status.STARTED, JobHandle.Status.CANCELLING);
+    when(jobExecutor.getJobExecutionInfo(startedJob.jobExecutionId()))
+        .thenReturn(
+            JobExecutionInfo.of(JobHandle.Status.STARTED),
+            JobExecutionInfo.of(JobHandle.Status.CANCELLING));
     stubEntityStoreUpdateToApply(startedJob);
 
     Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
@@ -1464,8 +1502,10 @@ public class TestJobManager {
     JobEntity queuedJob =
         newJobEntity("local-job-mine-2", JobHandle.Status.CANCELLING, 
Instant.now(), null);
     mockListActiveJobs(queuedJob);
-    when(jobExecutor.getJobStatus(queuedJob.jobExecutionId()))
-        .thenReturn(JobHandle.Status.QUEUED, JobHandle.Status.CANCELLED);
+    when(jobExecutor.getJobExecutionInfo(queuedJob.jobExecutionId()))
+        .thenReturn(
+            JobExecutionInfo.of(JobHandle.Status.QUEUED),
+            JobExecutionInfo.of(JobHandle.Status.CANCELLED));
     stubEntityStoreUpdateToApply(queuedJob);
 
     Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
@@ -1482,7 +1522,8 @@ public class TestJobManager {
         newJobEntity("local-job-mine-1", JobHandle.Status.CANCELLING, 
Instant.now(), null);
     mockListActiveJobs(job);
     when(jobExecutor.isJobStateNodeLocal()).thenReturn(true);
-    
when(jobExecutor.getJobStatus(job.jobExecutionId())).thenReturn(JobHandle.Status.STARTED);
+    when(jobExecutor.getJobExecutionInfo(job.jobExecutionId()))
+        .thenReturn(JobExecutionInfo.of(JobHandle.Status.STARTED));
     doThrow(new RuntimeException("cancel 
failed")).when(jobExecutor).cancelJob(any());
     stubEntityStoreUpdateToApply(job);
 
@@ -1502,14 +1543,16 @@ public class TestJobManager {
     JobEntity job =
         newJobEntity("external-job-1", JobHandle.Status.CANCELLING, 
Instant.now(), null);
     mockListActiveJobs(job);
-    
when(jobExecutor.getJobStatus(job.jobExecutionId())).thenReturn(JobHandle.Status.STARTED);
+    when(jobExecutor.getJobExecutionInfo(job.jobExecutionId()))
+        .thenReturn(JobExecutionInfo.of(JobHandle.Status.STARTED));
     stubEntityStoreUpdateToApply(job);
 
     Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
 
     verify(jobExecutor, never()).cancelJob(any());
-    // The observed STARTED status never regresses the CANCELLING job.
-    Assertions.assertEquals(JobHandle.Status.CANCELLING, 
captureUpdatedJobEntity(job).status());
+    // The observed STARTED status never regresses the CANCELLING job, so 
nothing is written.
+    verify(entityStore, never())
+        .update(any(), eq(JobEntity.class), eq(Entity.EntityType.JOB), any());
   }
 
   @Test
@@ -1519,7 +1562,8 @@ public class TestJobManager {
         newJobEntity("local-job-mine-1", JobHandle.Status.CANCELLING, 
Instant.now(), null);
     mockListActiveJobs(job);
     when(jobExecutor.isJobStateNodeLocal()).thenReturn(true);
-    
when(jobExecutor.getJobStatus(job.jobExecutionId())).thenReturn(JobHandle.Status.SUCCEEDED);
+    when(jobExecutor.getJobExecutionInfo(job.jobExecutionId()))
+        .thenReturn(JobExecutionInfo.of(JobHandle.Status.SUCCEEDED));
     stubEntityStoreUpdateToApply(job);
 
     Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
@@ -2234,10 +2278,10 @@ public class TestJobManager {
 
     when(jobManager.listJobs(metalake, Optional.empty()))
         .thenReturn(ImmutableList.of(conflictedJob, survivingJob));
-    when(jobExecutor.getJobStatus(conflictedJob.jobExecutionId()))
-        .thenReturn(JobHandle.Status.SUCCEEDED);
-    when(jobExecutor.getJobStatus(survivingJob.jobExecutionId()))
-        .thenReturn(JobHandle.Status.SUCCEEDED);
+    when(jobExecutor.getJobExecutionInfo(conflictedJob.jobExecutionId()))
+        .thenReturn(JobExecutionInfo.of(JobHandle.Status.SUCCEEDED));
+    when(jobExecutor.getJobExecutionInfo(survivingJob.jobExecutionId()))
+        .thenReturn(JobExecutionInfo.of(JobHandle.Status.SUCCEEDED));
 
     // A losing CAS must not stop this batch or future scheduled polls.
     NameIdentifier conflictedJobIdent = NameIdentifierUtil.ofJob(metalake, 
conflictedJob.name());
@@ -2316,6 +2360,317 @@ public class TestJobManager {
         .build();
   }
 
+  @Test
+  public void testPullJobStatusUsesExecutorReportedTimestamps() throws 
IOException {
+    // The job starts and finishes between two polls, so it is never observed 
as STARTED. The
+    // times reported by the job executor are recorded instead of the poll 
time.
+    Instant queuedAt = Instant.now().minusSeconds(60);
+    JobEntity job = newJobEntity("job-execution-1", JobHandle.Status.QUEUED, 
queuedAt, null);
+    mockListActiveJobs(job);
+    Instant startedAt = queuedAt.plusSeconds(1);
+    Instant finishedAt = queuedAt.plusSeconds(2);
+    doReturn(
+            JobExecutionInfo.builder()
+                .withStatus(JobHandle.Status.SUCCEEDED)
+                .withStartedAt(startedAt)
+                .withFinishedAt(finishedAt)
+                .build())
+        .when(jobExecutor)
+        .getJobExecutionInfo(job.jobExecutionId());
+    stubEntityStoreUpdateToApply(job);
+
+    Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
+
+    JobEntity updatedJob = captureUpdatedJobEntity(job);
+    Assertions.assertEquals(JobHandle.Status.SUCCEEDED, updatedJob.status());
+    Assertions.assertEquals(startedAt.toEpochMilli(), updatedJob.startedAt());
+    Assertions.assertEquals(finishedAt.toEpochMilli(), 
updatedJob.finishedAt());
+  }
+
+  @Test
+  public void testPullJobStatusReportedStartedAtReplacesPolledTime() throws 
IOException {
+    // An earlier poll recorded its own time as the started time, as the job 
executor didn't report
+    // one at that time. The time reported later replaces it.
+    Instant queuedAt = Instant.now().minusSeconds(60);
+    JobEntity job =
+        JobEntity.builder()
+            .withId(idGenerator.nextId())
+            .withJobExecutionId("job-execution-1")
+            .withNamespace(NamespaceUtil.ofJob(metalake))
+            .withJobTemplateName("shell_job")
+            .withStatus(JobHandle.Status.STARTED)
+            .withStartedAt(queuedAt.plusSeconds(30).toEpochMilli())
+            .withFinishedAt(0L)
+            
.withAuditInfo(AuditInfo.builder().withCreator("test").withCreateTime(queuedAt).build())
+            .build();
+    mockListActiveJobs(job);
+    Instant startedAt = queuedAt.plusSeconds(1);
+    doReturn(
+            JobExecutionInfo.builder()
+                .withStatus(JobHandle.Status.STARTED)
+                .withStartedAt(startedAt)
+                .build())
+        .when(jobExecutor)
+        .getJobExecutionInfo(job.jobExecutionId());
+    stubEntityStoreUpdateToApply(job);
+
+    Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
+
+    // The job is updated even though its status doesn't change.
+    JobEntity updatedJob = captureUpdatedJobEntity(job);
+    Assertions.assertEquals(JobHandle.Status.STARTED, updatedJob.status());
+    Assertions.assertEquals(startedAt.toEpochMilli(), updatedJob.startedAt());
+    Assertions.assertEquals(0L, updatedJob.finishedAt());
+  }
+
+  @Test
+  public void testPullJobStatusSkipsUpdateWhenNothingChanges() throws 
IOException {
+    Instant queuedAt = Instant.now().minusSeconds(60);
+    Instant startedAt = queuedAt.plusSeconds(1);
+    JobEntity job =
+        JobEntity.builder()
+            .withId(idGenerator.nextId())
+            .withJobExecutionId("job-execution-1")
+            .withNamespace(NamespaceUtil.ofJob(metalake))
+            .withJobTemplateName("shell_job")
+            .withStatus(JobHandle.Status.STARTED)
+            .withStartedAt(startedAt.toEpochMilli())
+            .withFinishedAt(0L)
+            
.withAuditInfo(AuditInfo.builder().withCreator("test").withCreateTime(queuedAt).build())
+            .build();
+    mockListActiveJobs(job);
+    doReturn(
+            JobExecutionInfo.builder()
+                .withStatus(JobHandle.Status.STARTED)
+                .withStartedAt(startedAt)
+                .build())
+        .when(jobExecutor)
+        .getJobExecutionInfo(job.jobExecutionId());
+
+    Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
+
+    verify(entityStore, never())
+        .update(any(), eq(JobEntity.class), eq(Entity.EntityType.JOB), any());
+  }
+
+  @Test
+  public void testPullJobStatusIgnoresInconsistentReportedTimestamps() throws 
IOException {
+    Instant queuedAt = Instant.now().minusSeconds(60);
+
+    // A queued job has no started time, so a reported one is ignored and 
nothing changes.
+    JobEntity queuedJob = newJobEntity("job-execution-1", 
JobHandle.Status.QUEUED, queuedAt, null);
+    mockListActiveJobs(queuedJob);
+    doReturn(
+            JobExecutionInfo.builder()
+                .withStatus(JobHandle.Status.QUEUED)
+                .withStartedAt(queuedAt.plusSeconds(1))
+                .build())
+        .when(jobExecutor)
+        .getJobExecutionInfo(queuedJob.jobExecutionId());
+    Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
+    verify(entityStore, never())
+        .update(any(), eq(JobEntity.class), eq(Entity.EntityType.JOB), any());
+
+    // A running job has no finished time, so a reported one is ignored.
+    JobEntity startingJob =
+        newJobEntity("job-execution-2", JobHandle.Status.QUEUED, queuedAt, 
null);
+    mockListActiveJobs(startingJob);
+    Instant startedAt = queuedAt.plusSeconds(1);
+    doReturn(
+            JobExecutionInfo.builder()
+                .withStatus(JobHandle.Status.STARTED)
+                .withStartedAt(startedAt)
+                .withFinishedAt(queuedAt.plusSeconds(2))
+                .build())
+        .when(jobExecutor)
+        .getJobExecutionInfo(startingJob.jobExecutionId());
+    stubEntityStoreUpdateToApply(startingJob);
+    Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
+
+    JobEntity startedJob = captureUpdatedJobEntity(startingJob);
+    Assertions.assertEquals(JobHandle.Status.STARTED, startedJob.status());
+    Assertions.assertEquals(startedAt.toEpochMilli(), startedJob.startedAt());
+    Assertions.assertEquals(0L, startedJob.finishedAt());
+  }
+
+  @Test
+  public void testPullJobStatusCorrectsReportedTimestamps() throws IOException 
{
+    Instant queuedAt = Instant.now().minusSeconds(60);
+
+    // The clock of the job runner is behind, so the job is reported to start 
and finish before it
+    // was queued. Both times are raised to the queued time.
+    JobEntity skewedJob = newJobEntity("job-execution-1", 
JobHandle.Status.QUEUED, queuedAt, null);
+    mockListActiveJobs(skewedJob);
+    doReturn(
+            JobExecutionInfo.builder()
+                .withStatus(JobHandle.Status.SUCCEEDED)
+                .withStartedAt(queuedAt.minusSeconds(2))
+                .withFinishedAt(queuedAt.minusSeconds(1))
+                .build())
+        .when(jobExecutor)
+        .getJobExecutionInfo(skewedJob.jobExecutionId());
+    stubEntityStoreUpdateToApply(skewedJob);
+    Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
+
+    JobEntity correctedJob = captureUpdatedJobEntity(skewedJob);
+    Assertions.assertEquals(queuedAt.toEpochMilli(), correctedJob.startedAt());
+    Assertions.assertEquals(queuedAt.toEpochMilli(), 
correctedJob.finishedAt());
+
+    // The job is reported to finish before it started. The started time is 
dropped, while the
+    // finished time is kept for the cleanup of the finished job.
+    Mockito.clearInvocations(entityStore);
+    JobEntity invalidJob = newJobEntity("job-execution-2", 
JobHandle.Status.QUEUED, queuedAt, null);
+    mockListActiveJobs(invalidJob);
+    Instant finishedAt = queuedAt.plusSeconds(1);
+    doReturn(
+            JobExecutionInfo.builder()
+                .withStatus(JobHandle.Status.FAILED)
+                .withStartedAt(queuedAt.plusSeconds(2))
+                .withFinishedAt(finishedAt)
+                .build())
+        .when(jobExecutor)
+        .getJobExecutionInfo(invalidJob.jobExecutionId());
+    stubEntityStoreUpdateToApply(invalidJob);
+    Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
+
+    JobEntity failedJob = captureUpdatedJobEntity(invalidJob);
+    Assertions.assertEquals(JobHandle.Status.FAILED, failedJob.status());
+    Assertions.assertEquals(0L, failedJob.startedAt());
+    Assertions.assertEquals(finishedAt.toEpochMilli(), failedJob.finishedAt());
+  }
+
+  @Test
+  public void 
testPullJobStatusDropsPolledStartedAtLaterThanReportedFinishedAt()
+      throws IOException {
+    // An earlier poll recorded its own time as the started time, and the job 
runner, whose clock
+    // is behind, reports the job finished before that. The recorded started 
time can't be right,
+    // so it is dropped, while the finished time is kept for the cleanup of 
the finished job.
+    Instant queuedAt = Instant.now().minusSeconds(60);
+    JobEntity job =
+        newStartedJobEntity("job-execution-1", queuedAt, 
queuedAt.plusSeconds(30).toEpochMilli());
+    mockListActiveJobs(job);
+    Instant finishedAt = queuedAt.plusSeconds(10);
+    when(jobExecutor.getJobExecutionInfo(job.jobExecutionId()))
+        .thenReturn(
+            JobExecutionInfo.builder()
+                .withStatus(JobHandle.Status.SUCCEEDED)
+                .withFinishedAt(finishedAt)
+                .build());
+    stubEntityStoreUpdateToApply(job);
+
+    Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
+
+    JobEntity updatedJob = captureUpdatedJobEntity(job);
+    Assertions.assertEquals(JobHandle.Status.SUCCEEDED, updatedJob.status());
+    Assertions.assertEquals(0L, updatedJob.startedAt());
+    Assertions.assertEquals(finishedAt.toEpochMilli(), 
updatedJob.finishedAt());
+  }
+
+  @Test
+  public void testPullJobStatusKeepsCancellingJobUntilItFinishes() throws 
IOException {
+    // A CANCELLING job is only updated once it finishes, even if the job 
executor reports a
+    // started time for it before that.
+    Instant queuedAt = Instant.now().minusSeconds(60);
+    JobEntity job = newJobEntity("external-job-1", 
JobHandle.Status.CANCELLING, queuedAt, null);
+    mockListActiveJobs(job);
+    Instant startedAt = queuedAt.plusSeconds(1);
+    when(jobExecutor.getJobExecutionInfo(job.jobExecutionId()))
+        .thenReturn(
+            JobExecutionInfo.builder()
+                .withStatus(JobHandle.Status.STARTED)
+                .withStartedAt(startedAt)
+                .build());
+
+    Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
+    verify(entityStore, never())
+        .update(any(), eq(JobEntity.class), eq(Entity.EntityType.JOB), any());
+
+    // Once the job is cancelled, it is updated with the started time as well.
+    Instant finishedAt = queuedAt.plusSeconds(2);
+    when(jobExecutor.getJobExecutionInfo(job.jobExecutionId()))
+        .thenReturn(
+            JobExecutionInfo.builder()
+                .withStatus(JobHandle.Status.CANCELLED)
+                .withStartedAt(startedAt)
+                .withFinishedAt(finishedAt)
+                .build());
+    stubEntityStoreUpdateToApply(job);
+
+    Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
+
+    JobEntity cancelledJob = captureUpdatedJobEntity(job);
+    Assertions.assertEquals(JobHandle.Status.CANCELLED, cancelledJob.status());
+    Assertions.assertEquals(startedAt.toEpochMilli(), 
cancelledJob.startedAt());
+    Assertions.assertEquals(finishedAt.toEpochMilli(), 
cancelledJob.finishedAt());
+  }
+
+  @Test
+  public void testPullJobStatusIgnoresUnusableReportedTimestamps() throws 
IOException {
+    // Times that don't fit in epoch milliseconds, or are far in the future, 
are dropped instead of
+    // failing the status pull, and the time of the pull is used for the 
finished job instead.
+    Instant queuedAt = Instant.now().minusSeconds(60);
+    JobEntity job = newJobEntity("job-execution-1", JobHandle.Status.QUEUED, 
queuedAt, null);
+    mockListActiveJobs(job);
+    when(jobExecutor.getJobExecutionInfo(job.jobExecutionId()))
+        .thenReturn(
+            JobExecutionInfo.builder()
+                .withStatus(JobHandle.Status.SUCCEEDED)
+                .withStartedAt(Instant.now().plus(Duration.ofDays(365)))
+                .withFinishedAt(Instant.MAX)
+                .build());
+    stubEntityStoreUpdateToApply(job);
+
+    Instant pulledAfter = Instant.now();
+    Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
+
+    JobEntity updatedJob = captureUpdatedJobEntity(job);
+    Assertions.assertEquals(JobHandle.Status.SUCCEEDED, updatedJob.status());
+    Assertions.assertEquals(0L, updatedJob.startedAt());
+    Assertions.assertTrue(updatedJob.finishedAt() >= 
pulledAfter.toEpochMilli());
+    Assertions.assertTrue(updatedJob.finishedAt() <= 
Instant.now().toEpochMilli());
+  }
+
+  @Test
+  public void testPullJobStatusContinuesAfterFailingToUpdateAJob() throws 
IOException {
+    // A failure on one job must not stop the status pull of the others, or 
escape the scheduled
+    // task, which would cancel all its later runs.
+    Instant queuedAt = Instant.now().minusSeconds(60);
+    JobEntity failingJob = newJobEntity("job-execution-1", 
JobHandle.Status.QUEUED, queuedAt, null);
+    JobEntity job = newJobEntity("job-execution-2", JobHandle.Status.QUEUED, 
queuedAt, null);
+    mockListActiveJobs(failingJob, job);
+    when(jobExecutor.getJobExecutionInfo(any()))
+        .thenReturn(JobExecutionInfo.of(JobHandle.Status.SUCCEEDED));
+    when(entityStore.update(
+            eq(NameIdentifierUtil.ofJob(metalake, failingJob.name())),
+            eq(JobEntity.class),
+            eq(Entity.EntityType.JOB),
+            any()))
+        .thenThrow(new RuntimeException("update failed"));
+    stubEntityStoreUpdateToApply(job, job);
+
+    Assertions.assertDoesNotThrow(() -> jobManager.pullAndUpdateJobStatus());
+
+    verify(entityStore)
+        .update(
+            eq(NameIdentifierUtil.ofJob(metalake, job.name())),
+            eq(JobEntity.class),
+            eq(Entity.EntityType.JOB),
+            any());
+  }
+
+  private JobEntity newStartedJobEntity(String executionId, Instant queuedAt, 
long startedAt) {
+    return JobEntity.builder()
+        .withId(idGenerator.nextId())
+        .withJobExecutionId(executionId)
+        .withNamespace(NamespaceUtil.ofJob(metalake))
+        .withJobTemplateName("shell_job")
+        .withStatus(JobHandle.Status.STARTED)
+        .withStartedAt(startedAt)
+        .withFinishedAt(0L)
+        
.withAuditInfo(AuditInfo.builder().withCreator("test").withCreateTime(queuedAt).build())
+        .build();
+  }
+
   private JobEntity newJobEntity(String templateName, JobHandle.Status status) 
{
     Random rand = new Random();
     return JobEntity.builder()
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 f2f50cda2d..487645a10d 100644
--- a/core/src/test/java/org/apache/gravitino/job/TestJobManagerMultiNode.java
+++ b/core/src/test/java/org/apache/gravitino/job/TestJobManagerMultiNode.java
@@ -40,6 +40,7 @@ import org.apache.gravitino.Entity;
 import org.apache.gravitino.EntityStore;
 import org.apache.gravitino.GravitinoEnv;
 import org.apache.gravitino.cache.NoOpsCache;
+import org.apache.gravitino.connector.job.JobExecutionInfo;
 import org.apache.gravitino.connector.job.JobExecutor;
 import org.apache.gravitino.exceptions.NoSuchJobException;
 import org.apache.gravitino.job.local.LocalJobExecutor;
@@ -146,6 +147,28 @@ public class TestJobManagerMultiNode extends 
TestJDBCBackend {
     Assertions.assertTrue(afterNodeAPull.finishedAt() > 0);
   }
 
+  @TestTemplate
+  public void testShortJobRecordsActualTimestamps() throws IOException {
+    // The job starts and finishes between two polls, so it is never observed 
as STARTED.
+    JobEntity job = nodeA.runJob(METALAKE, TEMPLATE, 
ImmutableMap.of("seconds", "1"));
+    Awaitility.await()
+        .atMost(1, TimeUnit.MINUTES)
+        .until(() -> executorA.getJobStatus(job.jobExecutionId()) == 
JobHandle.Status.SUCCEEDED);
+    JobExecutionInfo info = 
executorA.getJobExecutionInfo(job.jobExecutionId());
+
+    // The job is polled after it finished, the recorded times are still when 
it actually ran.
+    nodeA.pullAndUpdateJobStatus();
+
+    JobEntity finished = getJob(job.name());
+    Assertions.assertEquals(JobHandle.Status.SUCCEEDED, finished.status());
+    // The metadata store keeps the timestamps in milliseconds.
+    Assertions.assertEquals(info.startedAt().toEpochMilli(), 
finished.startedAt());
+    Assertions.assertEquals(info.finishedAt().toEpochMilli(), 
finished.finishedAt());
+    Assertions.assertTrue(
+        finished.startedAt() >= job.auditInfo().createTime().toEpochMilli(), 
finished.toString());
+    Assertions.assertTrue(finished.finishedAt() - finished.startedAt() >= 900, 
finished.toString());
+  }
+
   @TestTemplate
   public void testCancelJobFromAnotherNode() throws IOException {
     JobEntity job = runLongJobOnNodeA();
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 ee2d0dae25..7908021347 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
@@ -30,6 +30,8 @@ import java.net.URL;
 import java.nio.charset.StandardCharsets;
 import java.nio.file.Files;
 import java.nio.file.attribute.FileTime;
+import java.time.Duration;
+import java.time.Instant;
 import java.util.Collections;
 import java.util.List;
 import java.util.Map;
@@ -41,6 +43,7 @@ import java.util.concurrent.Future;
 import java.util.concurrent.TimeUnit;
 import org.apache.commons.io.FileUtils;
 import org.apache.commons.lang3.StringUtils;
+import org.apache.gravitino.connector.job.JobExecutionInfo;
 import org.apache.gravitino.connector.job.JobExecutor;
 import org.apache.gravitino.exceptions.NoSuchJobException;
 import org.apache.gravitino.job.JobHandle;
@@ -1234,13 +1237,111 @@ public class TestLocalJobExecutor {
     }
   }
 
+  @Test
+  public void testJobExecutionInfoOfFinishedJobs() throws IOException {
+    Instant submittedAt = Instant.now();
+    String succeededJobId =
+        jobExecutor.submitJob(newScriptJobTemplate("succeed", "sleep 1\nexit 
0"));
+    String failedJobId = jobExecutor.submitJob(newScriptJobTemplate("fail", 
"sleep 1\nexit 1"));
+
+    Awaitility.await()
+        .atMost(1, TimeUnit.MINUTES)
+        .until(
+            () ->
+                jobExecutor.getJobStatus(succeededJobId) == 
JobHandle.Status.SUCCEEDED
+                    && jobExecutor.getJobStatus(failedJobId) == 
JobHandle.Status.FAILED);
+
+    // The jobs are never polled while they run, but their snapshots still 
carry when they
+    // actually started and finished.
+    for (String jobId : Lists.newArrayList(succeededJobId, failedJobId)) {
+      JobExecutionInfo info = jobExecutor.getJobExecutionInfo(jobId);
+      Assertions.assertNotNull(info.startedAt(), jobId);
+      Assertions.assertNotNull(info.finishedAt(), jobId);
+      Assertions.assertFalse(info.startedAt().isBefore(submittedAt), jobId);
+      Assertions.assertTrue(
+          Duration.between(info.startedAt(), info.finishedAt()).toMillis() >= 
900, info.toString());
+    }
+  }
+
+  @Test
+  public void testJobExecutionInfoOfCancelledJobs() throws IOException {
+    LocalJobExecutor executor = new LocalJobExecutor();
+    try {
+      executor.initialize(
+          
withStagingDir(ImmutableMap.of(LocalJobExecutorConfigs.MAX_RUNNING_JOBS, "1")));
+      String runningJobId = executor.submitJob(newSleepJobTemplate("sleep"));
+      Awaitility.await()
+          .atMost(1, TimeUnit.MINUTES)
+          .until(() -> executor.getJobStatus(runningJobId) == 
JobHandle.Status.STARTED);
+      // The other job waits in the queue, as only one job can run at a time.
+      String queuedJobId = executor.submitJob(newScriptJobTemplate("queued", 
"exit 0"));
+
+      JobExecutionInfo started = executor.getJobExecutionInfo(runningJobId);
+      Assertions.assertNotNull(started.startedAt());
+      Assertions.assertNull(started.finishedAt());
+      Assertions.assertEquals(
+          JobExecutionInfo.of(JobHandle.Status.QUEUED), 
executor.getJobExecutionInfo(queuedJobId));
+
+      // A job cancelled from the queue never started.
+      executor.cancelJob(queuedJobId);
+      JobExecutionInfo cancelledQueued = 
executor.getJobExecutionInfo(queuedJobId);
+      Assertions.assertEquals(JobHandle.Status.CANCELLED, 
cancelledQueued.status());
+      Assertions.assertNull(cancelledQueued.startedAt());
+      Assertions.assertNotNull(cancelledQueued.finishedAt());
+
+      // A running job keeps its started time while it is being cancelled and 
after that.
+      executor.cancelJob(runningJobId);
+      Awaitility.await()
+          .atMost(1, TimeUnit.MINUTES)
+          .until(() -> executor.getJobStatus(runningJobId) == 
JobHandle.Status.CANCELLED);
+      JobExecutionInfo cancelled = executor.getJobExecutionInfo(runningJobId);
+      Assertions.assertEquals(started.startedAt(), cancelled.startedAt());
+      Assertions.assertNotNull(cancelled.finishedAt());
+      
Assertions.assertFalse(cancelled.finishedAt().isBefore(cancelled.startedAt()));
+    } finally {
+      executor.close();
+    }
+  }
+
+  @Test
+  public void testCleanupOnlyRemovesFinishedJobs() throws IOException {
+    LocalJobExecutor executor = new LocalJobExecutor();
+    try {
+      executor.initialize(
+          withStagingDir(
+              ImmutableMap.of(
+                  LocalJobExecutorConfigs.JOB_STATUS_KEEP_TIME_MS,
+                  "10",
+                  LocalJobExecutorConfigs.MAX_RUNNING_JOBS,
+                  "2")));
+      String runningJobId = executor.submitJob(newSleepJobTemplate("sleep"));
+      Awaitility.await()
+          .atMost(1, TimeUnit.MINUTES)
+          .until(() -> executor.getJobStatus(runningJobId) == 
JobHandle.Status.STARTED);
+      String finishedJobId = executor.submitJob(newScriptJobTemplate("finish", 
"exit 0"));
+
+      // The finished job is removed once it has been kept for the keep time, 
while the running job
+      // is kept however long it runs.
+      Awaitility.await()
+          .atMost(1, TimeUnit.MINUTES)
+          .until(() -> !hasJobStatus(executor, finishedJobId));
+      Assertions.assertEquals(JobHandle.Status.STARTED, 
executor.getJobStatus(runningJobId));
+    } finally {
+      executor.close();
+    }
+  }
+
   private JobTemplate newSleepJobTemplate(String name) throws IOException {
+    // Exec the sleep, so that killing the job process also stops the sleep.
+    return newScriptJobTemplate(name, "exec sleep 600");
+  }
+
+  private JobTemplate newScriptJobTemplate(String name, String commands) 
throws IOException {
     // The job runs in the directory of its executable, so give each job its 
own directory.
     File jobDir = new File(workingDir, name);
     Assertions.assertTrue(jobDir.mkdirs());
-    File script = new File(jobDir, "sleep.sh");
-    // Exec the sleep, so that killing the job process also stops the sleep.
-    Files.writeString(script.toPath(), "#!/bin/bash\nexec sleep 600\n");
+    File script = new File(jobDir, "job.sh");
+    Files.writeString(script.toPath(), "#!/bin/bash\n" + commands + "\n");
     Assertions.assertTrue(script.setExecutable(true));
 
     return ShellJobTemplate.builder()
diff --git a/docs/development/custom-job-executor.md 
b/docs/development/custom-job-executor.md
index f706422147..28efdee10b 100644
--- a/docs/development/custom-job-executor.md
+++ b/docs/development/custom-job-executor.md
@@ -17,6 +17,13 @@ Gravitino's job system is extensible: you can implement your 
own job executor
 to run jobs in a distributed environment. Refer to the interface `JobExecutor` 
in the
 code 
[here](https://github.com/apache/gravitino/blob/main/core/src/main/java/org/apache/gravitino/connector/job/JobExecutor.java).
 
+Gravitino tracks the jobs by pulling their state with `getJobExecutionInfo`, 
which every job executor
+must implement. It returns the job's status, and when the job actually started 
and finished. Report
+these times whenever the job runner provides them: Gravitino pulls job states 
only every
+`gravitino.job.statusPullIntervalInMs`, so without them it can only record 
when it observed each
+status change, and a job that starts and finishes between two pulls gets no 
start time. Once a job
+has started, keep reporting its start time in every later state, including the 
finished one.
+
 After you implement your own job executor, you need to register it in the 
Gravitino server by
 using the `gravitino.conf` file. For example, if you have implemented a job 
executor named
 `airflow`, you need to configure it as follows:
diff --git a/docs/manage-jobs-in-gravitino.md b/docs/manage-jobs-in-gravitino.md
index 019e273ea3..165c660b59 100644
--- a/docs/manage-jobs-in-gravitino.md
+++ b/docs/manage-jobs-in-gravitino.md
@@ -224,6 +224,28 @@ cancelling = client.cancel_job(job_id)
 Cancelling is a request rather than an instant. The job moves to `CANCELLING` 
and then to
 `CANCELLED`, and one that finishes first keeps the status it finished with.
 
+A job carries three timestamps:
+
+- `queuedAt`: when Gravitino submitted the job to the job executor.
+- `startedAt`: when the job started executing.
+- `finishedAt`: when the job finished.
+
+Gravitino pulls job statuses from the job executor every 
`gravitino.job.statusPullIntervalInMs`,
+so a job's status can lag behind by up to this interval. The timestamps 
usually don't lag: the local
+job executor reports when each job actually started and finished, even for a 
job that starts and
+finishes between two pulls. A job executor that doesn't report these times 
gets the time Gravitino
+first observes the job running or finished instead. In that case, a job that 
finishes between two
+pulls has no `startedAt`.
+
+The local job executor only keeps a job's state in memory. When that state is 
lost before Gravitino
+records the finished job, the actual times are lost too, and the job's 
`finishedAt` is the time
+Gravitino marks it as finished instead:
+
+- If Gravitino can't pull the job within 
`gravitino.jobExecutor.local.jobStatusKeepTimeInMs` after it
+  finished, the job is marked as `FAILED`, or `CANCELLED` if it was being 
cancelled, on the next pull.
+- If the server running the job exits, the job is expired as described in
+  [Configurations for Local Job 
Executor](#configurations-for-local-job-executor).
+
 ### Get a Job's Output
 
 A job's captured stdout/stderr can be fetched alongside its metadata by asking 
for it explicitly.
diff --git a/docs/open-api/jobs.yaml b/docs/open-api/jobs.yaml
index 856a999fb1..36a63aa38e 100644
--- a/docs/open-api/jobs.yaml
+++ b/docs/open-api/jobs.yaml
@@ -645,7 +645,7 @@ components:
           type: string
           format: date-time
           nullable: true
-          description: The time when the job started execution, or null if the 
job has not started execution yet
+          description: The time when the job started execution, or null if the 
job has not started execution yet. A finished job may also have no started 
time, if the job executor doesn't report when the job started and Gravitino 
didn't observe the job running
         finishedAt:
           type: string
           format: date-time

Reply via email to