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

gyfora pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/flink-kubernetes-operator.git


The following commit(s) were added to refs/heads/main by this push:
     new 492e14d1 [FLINK-40164] Complete the suspend when the cancellation is 
observed (#1168)
492e14d1 is described below

commit 492e14d137b882f99f3f769c1bd6ee3a6fcbfaaf
Author: Milind L <[email protected]>
AuthorDate: Wed Aug 5 13:00:15 2026 +0530

    [FLINK-40164] Complete the suspend when the cancellation is observed (#1168)
    
    Suspending a job through an asynchronous cancellation left the resource 
reporting the UPGRADING lifecycle state forever. An example that reproduces 
this exists in the ticket.
    
    I experienced it with FlinkSessionJobs (stateless or last-state upgrade 
mode), but even FlinkDeployments in application mode will suffer this (if 
`kubernetes.operator.job.upgrade.last-state.job-cancel.enabled` and 
`upgradeMode` is `last-state` for example).
    
    Savepoint upgrades complete synchronously and are not affected.
    
    To solve: In the observer, after marking resource suspended, we check if it 
was a user has set the desired state of the job to SUSPENDED and whether it's 
UPGRADING, and if so, mark the spec as stable and set 
`ReconciliationState.DEPLOYED`.
---
 .../operator/observer/JobStatusObserver.java       |  24 +++
 .../controller/FlinkDeploymentControllerTest.java  |  59 +++++++
 .../controller/FlinkSessionJobControllerTest.java  | 160 +++++++++++++++++
 .../operator/observer/JobStatusObserverTest.java   | 194 +++++++++++++++++++++
 4 files changed, 437 insertions(+)

diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/observer/JobStatusObserver.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/observer/JobStatusObserver.java
index 0c875613..e205e3b8 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/observer/JobStatusObserver.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/observer/JobStatusObserver.java
@@ -350,6 +350,7 @@ public class JobStatusObserver<R extends 
AbstractFlinkResource<?, ?>> {
                             && 
JobStatus.CANCELLING.equals(previousJobStatus))) {
                 // The job was cancelled
                 markSuspended(resource);
+                finalizeSuspendedUpgrade(resource);
             }
 
             recordJobErrorIfPresent(ctx, clusterJobStatus);
@@ -373,6 +374,29 @@ public class JobStatusObserver<R extends 
AbstractFlinkResource<?, ?>> {
                 });
     }
 
+    /**
+     * When a suspend is executed through an asynchronous cancellation the 
reconciler exits early
+     * and leaves the resource in the UPGRADING state until the cancellation 
completes. Once the job
+     * has ended there is nothing left to do for a suspend request, so we 
complete the upgrade here,
+     * otherwise the resource would report the UPGRADING lifecycle state 
indefinitely.
+     *
+     * <p>Must be called right after {@link 
#markSuspended(AbstractFlinkResource)}, which records
+     * the suspended job state that this completion makes stable.
+     *
+     * @param resource The Flink resource whose cancellation was just observed.
+     */
+    private static void finalizeSuspendedUpgrade(AbstractFlinkResource<?, ?> 
resource) {
+        var reconciliationStatus = 
resource.getStatus().getReconciliationStatus();
+        if (reconciliationStatus.getState() != ReconciliationState.UPGRADING
+                || resource.getSpec().getJob().getState() != 
JobState.SUSPENDED) {
+            return;
+        }
+        LOG.debug("Suspend completed after asynchronous cancellation");
+        reconciliationStatus.setState(ReconciliationState.DEPLOYED);
+        // Suspended specs are always marked stable, same as for synchronous 
suspends
+        reconciliationStatus.markReconciledSpecAsStable();
+    }
+
     private void recordJobErrorIfPresent(
             FlinkResourceContext<R> ctx, JobStatusMessage clusterJobStatus) {
         if (clusterJobStatus.getJobState() == JobStatus.FAILED) {
diff --git 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkDeploymentControllerTest.java
 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkDeploymentControllerTest.java
index 65cdfaf0..7581fce6 100644
--- 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkDeploymentControllerTest.java
+++ 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkDeploymentControllerTest.java
@@ -32,6 +32,7 @@ import 
org.apache.flink.kubernetes.operator.api.spec.UpgradeMode;
 import 
org.apache.flink.kubernetes.operator.api.status.FlinkDeploymentReconciliationStatus;
 import 
org.apache.flink.kubernetes.operator.api.status.JobManagerDeploymentStatus;
 import org.apache.flink.kubernetes.operator.api.status.JobStatus;
+import org.apache.flink.kubernetes.operator.api.status.ReconciliationState;
 import org.apache.flink.kubernetes.operator.api.status.ReconciliationStatus;
 import org.apache.flink.kubernetes.operator.api.status.TaskManagerInfo;
 import org.apache.flink.kubernetes.operator.config.FlinkConfigManager;
@@ -65,6 +66,7 @@ import java.util.Optional;
 import java.util.stream.Collectors;
 
 import static 
org.apache.flink.kubernetes.operator.TestUtils.MAX_RECONCILE_TIMES;
+import static 
org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions.OPERATOR_JOB_UPGRADE_LAST_STATE_CANCEL_JOB;
 import static 
org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions.OPERATOR_JOB_UPGRADE_LAST_STATE_FALLBACK_ENABLED;
 import static 
org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions.SNAPSHOT_RESOURCE_ENABLED;
 import static 
org.apache.flink.kubernetes.operator.utils.EventRecorder.Reason.ValidationError;
@@ -300,6 +302,63 @@ public class FlinkDeploymentControllerTest {
                 updateControl.getScheduleDelay().get());
     }
 
+    @Test
+    public void verifySuspendCompletesAfterAsyncCancellation() throws 
Exception {
+        FlinkDeployment appCluster = TestUtils.buildApplicationCluster();
+        appCluster.getSpec().getJob().setUpgradeMode(UpgradeMode.LAST_STATE);
+        appCluster
+                .getSpec()
+                .getFlinkConfiguration()
+                .put(OPERATOR_JOB_UPGRADE_LAST_STATE_CANCEL_JOB.key(), "true");
+
+        // Deploying an application cluster takes several reconciliations to 
reach a running job
+        for (int i = 0; i < 4; i++) {
+            testController.reconcile(appCluster, context);
+        }
+        assertEquals(
+                org.apache.flink.api.common.JobStatus.RUNNING,
+                appCluster.getStatus().getJobStatus().getState());
+
+        // Cancelling the job is the only asynchronous application mode 
suspend, the reconciler
+        // records UPGRADING and has to wait for the observer to see the 
cancellation complete
+        appCluster.getSpec().getJob().setState(JobState.SUSPENDED);
+        testController.reconcile(appCluster, context);
+        assertEquals(
+                org.apache.flink.api.common.JobStatus.CANCELLING,
+                appCluster.getStatus().getJobStatus().getState());
+        assertEquals(
+                ReconciliationState.UPGRADING,
+                appCluster.getStatus().getReconciliationStatus().getState());
+        assertEquals(ResourceLifecycleState.UPGRADING, 
appCluster.getStatus().getLifecycleState());
+
+        // Once the cancellation was observed there is nothing left to do for 
the suspend
+        testController.reconcile(appCluster, context);
+        assertSuspendCompleted(appCluster);
+
+        // Further reconciliations must not disturb the suspended state
+        for (int i = 0; i < MAX_RECONCILE_TIMES; i++) {
+            testController.reconcile(appCluster, context);
+            assertSuspendCompleted(appCluster);
+        }
+    }
+
+    private static void assertSuspendCompleted(FlinkDeployment appCluster) {
+        var status = appCluster.getStatus();
+        var reconciliationStatus = status.getReconciliationStatus();
+        var lastReconciledJobSpec = 
reconciliationStatus.deserializeLastReconciledSpec().getJob();
+
+        assertEquals(
+                org.apache.flink.api.common.JobStatus.CANCELED, 
status.getJobStatus().getState());
+        assertEquals(ReconciliationState.DEPLOYED, 
reconciliationStatus.getState());
+        assertEquals(JobState.SUSPENDED, lastReconciledJobSpec.getState());
+        // The upgrade mode recorded while cancelling describes how the state 
was preserved, it has
+        // to be kept so that the job can be restored later
+        assertEquals(UpgradeMode.SAVEPOINT, 
lastReconciledJobSpec.getUpgradeMode());
+        assertTrue(reconciliationStatus.isLastReconciledSpecStable());
+        assertEquals(ResourceLifecycleState.SUSPENDED, 
status.getLifecycleState());
+        assertNull(status.getError());
+    }
+
     @ParameterizedTest
     
@MethodSource("org.apache.flink.kubernetes.operator.TestUtils#flinkVersions")
     public void verifyUpgradeFromSavepointLegacyMode(FlinkVersion 
flinkVersion) throws Exception {
diff --git 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkSessionJobControllerTest.java
 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkSessionJobControllerTest.java
index af4302be..5593cfbc 100644
--- 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkSessionJobControllerTest.java
+++ 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/controller/FlinkSessionJobControllerTest.java
@@ -54,6 +54,7 @@ import java.util.Map;
 import java.util.Optional;
 import java.util.stream.Collectors;
 
+import static org.apache.flink.api.common.JobStatus.CANCELED;
 import static org.apache.flink.api.common.JobStatus.CANCELLING;
 import static org.apache.flink.api.common.JobStatus.FINISHED;
 import static org.apache.flink.api.common.JobStatus.RECONCILING;
@@ -328,6 +329,150 @@ class FlinkSessionJobControllerTest {
         assertEquals("cp2", 
sessionJob.getStatus().getJobStatus().getUpgradeSavepointPath());
     }
 
+    @Test
+    public void verifySuspendCompletesAfterAsyncCancellation() throws 
Exception {
+        testController.reconcile(sessionJob, context);
+        testController.reconcile(sessionJob, context);
+        assertEquals(RUNNING, 
sessionJob.getStatus().getJobStatus().getState());
+
+        // Suspending a stateless session job is executed through an 
asynchronous cancellation, the
+        // reconciler can only initiate it and has to wait for the observer to 
see it complete
+        sessionJob.getSpec().getJob().setState(JobState.SUSPENDED);
+        testController.reconcile(sessionJob, context);
+        assertEquals(CANCELLING, 
sessionJob.getStatus().getJobStatus().getState());
+        assertEquals(
+                ReconciliationState.UPGRADING,
+                sessionJob.getStatus().getReconciliationStatus().getState());
+        assertEquals(ResourceLifecycleState.UPGRADING, 
sessionJob.getStatus().getLifecycleState());
+
+        // Once the cancellation was observed there is nothing left to do for 
the suspend
+        testController.reconcile(sessionJob, context);
+        assertSuspendCompleted(sessionJob, UpgradeMode.STATELESS);
+
+        // Further reconciliations must not disturb the suspended state
+        for (int i = 0; i < MAX_RECONCILE_TIMES; i++) {
+            testController.reconcile(sessionJob, context);
+            assertSuspendCompleted(sessionJob, UpgradeMode.STATELESS);
+        }
+
+        // The job can be resumed
+        flinkService.clearJobsInTerminalState();
+        sessionJob.getSpec().getJob().setState(JobState.RUNNING);
+        testController.reconcile(sessionJob, context);
+        assertEquals(1, flinkService.listJobs().size());
+        assertEquals(RECONCILING, 
sessionJob.getStatus().getJobStatus().getState());
+        assertEquals(
+                ReconciliationState.DEPLOYED,
+                sessionJob.getStatus().getReconciliationStatus().getState());
+    }
+
+    @ParameterizedTest
+    @EnumSource(
+            value = org.apache.flink.api.common.JobStatus.class,
+            mode = EnumSource.Mode.INCLUDE,
+            names = {"FINISHED", "FAILED"})
+    public void verifySuspendCompletesWhenJobEndsWhileCancelling(
+            org.apache.flink.api.common.JobStatus terminalStatus) throws 
Exception {
+        testController.reconcile(sessionJob, context);
+        testController.reconcile(sessionJob, context);
+        assertEquals(RUNNING, 
sessionJob.getStatus().getJobStatus().getState());
+
+        sessionJob.getSpec().getJob().setState(JobState.SUSPENDED);
+        testController.reconcile(sessionJob, context);
+        assertEquals(CANCELLING, 
sessionJob.getStatus().getJobStatus().getState());
+
+        // The job reaches a globally terminal state on its own before the 
cancellation completes
+        var job = flinkService.listJobs().get(0);
+        job.f1 =
+                new JobStatusMessage(
+                        job.f1.getJobId(),
+                        job.f1.getJobName(),
+                        terminalStatus,
+                        job.f1.getStartTime());
+
+        testController.reconcile(sessionJob, context);
+
+        // The suspend is over however the job ended, so the resource must be 
reported as suspended
+        // instead of waiting for a cancellation that will never be observed
+        var status = sessionJob.getStatus();
+        var reconciliationStatus = status.getReconciliationStatus();
+        assertEquals(terminalStatus, status.getJobStatus().getState());
+        assertEquals(ReconciliationState.DEPLOYED, 
reconciliationStatus.getState());
+        assertEquals(
+                JobState.SUSPENDED,
+                
reconciliationStatus.deserializeLastReconciledSpec().getJob().getState());
+        assertTrue(reconciliationStatus.isLastReconciledSpecStable());
+        assertEquals(ResourceLifecycleState.SUSPENDED, 
status.getLifecycleState());
+        // The user asked for a suspend, the terminated job must not be 
submitted again
+        assertEquals(1, flinkService.listJobs().size());
+        assertEquals(0, flinkService.getRunningCount());
+    }
+
+    @Test
+    public void verifyLastStateSuspendCompletesAfterAsyncCancellation() throws 
Exception {
+        sessionJob.getSpec().getJob().setUpgradeMode(UpgradeMode.LAST_STATE);
+        testController.reconcile(sessionJob, context);
+        testController.reconcile(sessionJob, context);
+        assertEquals(RUNNING, 
sessionJob.getStatus().getJobStatus().getState());
+
+        // Simulate completed checkpoints
+        flinkService.setCheckpointInfo(
+                Tuple2.of(
+                        Optional.of(
+                                new 
CheckpointHistoryWrapper.CompletedCheckpointInfo(
+                                        0, "cp1", System.currentTimeMillis())),
+                        Optional.empty()));
+
+        sessionJob.getSpec().getJob().setState(JobState.SUSPENDED);
+        testController.reconcile(sessionJob, context);
+        assertEquals(CANCELLING, 
sessionJob.getStatus().getJobStatus().getState());
+
+        testController.reconcile(sessionJob, context);
+        // The upgrade mode recorded while cancelling describes how the state 
was preserved, it must
+        // be kept so that the job can be restored from the checkpoint later
+        assertSuspendCompleted(sessionJob, UpgradeMode.SAVEPOINT);
+        assertEquals("cp1", 
sessionJob.getStatus().getJobStatus().getUpgradeSavepointPath());
+
+        // Resuming restores from the checkpoint taken during the cancellation
+        flinkService.clearJobsInTerminalState();
+        sessionJob.getSpec().getJob().setState(JobState.RUNNING);
+        testController.reconcile(sessionJob, context);
+        var jobs = flinkService.listJobs();
+        assertEquals(1, jobs.size());
+        assertEquals("cp1", jobs.get(0).f0);
+    }
+
+    @Test
+    public void verifyExternallyCancelledJobIsRestarted() throws Exception {
+        testController.reconcile(sessionJob, context);
+        testController.reconcile(sessionJob, context);
+        assertEquals(RUNNING, 
sessionJob.getStatus().getJobStatus().getState());
+
+        // The job is cancelled outside the operator while the target state is 
still RUNNING
+        var job = flinkService.listJobs().get(0);
+        job.f1 =
+                new JobStatusMessage(
+                        job.f1.getJobId(), job.f1.getJobName(), CANCELED, 
job.f1.getStartTime());
+
+        testController.reconcile(sessionJob, context);
+
+        // The operator must resubmit the job instead of reporting the 
resource as suspended, so the
+        // cancelled job is joined by a freshly submitted one
+        assertEquals(2, flinkService.listJobs().size());
+        assertEquals(1, flinkService.getRunningCount());
+        assertEquals(
+                ReconciliationState.DEPLOYED,
+                sessionJob.getStatus().getReconciliationStatus().getState());
+        assertEquals(
+                JobState.RUNNING,
+                sessionJob
+                        .getStatus()
+                        .getReconciliationStatus()
+                        .deserializeLastReconciledSpec()
+                        .getJob()
+                        .getState());
+    }
+
     @Test
     public void verifyLastStateUpgradeFailure() throws Exception {
         sessionJob.getSpec().getJob().setUpgradeMode(UpgradeMode.LAST_STATE);
@@ -895,4 +1040,19 @@ class FlinkSessionJobControllerTest {
         assertEquals("CheckpointNotFound", errorEvent.get().getReason());
         assertTrue(errorEvent.get().getMessage().contains("not externally 
addressable"));
     }
+
+    private static void assertSuspendCompleted(
+            FlinkSessionJob sessionJob, UpgradeMode expectedUpgradeMode) {
+        var status = sessionJob.getStatus();
+        var reconciliationStatus = status.getReconciliationStatus();
+        var lastReconciledJobSpec = 
reconciliationStatus.deserializeLastReconciledSpec().getJob();
+
+        assertEquals(CANCELED, status.getJobStatus().getState());
+        assertEquals(ReconciliationState.DEPLOYED, 
reconciliationStatus.getState());
+        assertEquals(JobState.SUSPENDED, lastReconciledJobSpec.getState());
+        assertEquals(expectedUpgradeMode, 
lastReconciledJobSpec.getUpgradeMode());
+        assertTrue(reconciliationStatus.isLastReconciledSpecStable());
+        assertEquals(ResourceLifecycleState.SUSPENDED, 
status.getLifecycleState());
+        assertNull(status.getError());
+    }
 }
diff --git 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/observer/JobStatusObserverTest.java
 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/observer/JobStatusObserverTest.java
index f3ff7a8b..67a0a929 100644
--- 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/observer/JobStatusObserverTest.java
+++ 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/observer/JobStatusObserverTest.java
@@ -19,6 +19,7 @@ package org.apache.flink.kubernetes.operator.observer;
 
 import org.apache.flink.api.common.JobID;
 import org.apache.flink.api.common.JobStatus;
+import org.apache.flink.api.java.tuple.Tuple3;
 import org.apache.flink.configuration.Configuration;
 import org.apache.flink.configuration.PipelineOptionsInternal;
 import org.apache.flink.kubernetes.operator.OperatorTestBase;
@@ -26,12 +27,15 @@ import org.apache.flink.kubernetes.operator.TestUtils;
 import org.apache.flink.kubernetes.operator.api.AbstractFlinkResource;
 import org.apache.flink.kubernetes.operator.api.FlinkDeployment;
 import org.apache.flink.kubernetes.operator.api.FlinkSessionJob;
+import 
org.apache.flink.kubernetes.operator.api.lifecycle.ResourceLifecycleState;
 import org.apache.flink.kubernetes.operator.api.spec.JobState;
 import org.apache.flink.kubernetes.operator.api.spec.UpgradeMode;
+import org.apache.flink.kubernetes.operator.api.status.ReconciliationState;
 import 
org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions;
 import org.apache.flink.kubernetes.operator.controller.FlinkResourceContext;
 import org.apache.flink.kubernetes.operator.reconciler.ReconciliationUtils;
 import org.apache.flink.kubernetes.operator.utils.EventRecorder;
+import org.apache.flink.runtime.client.JobStatusMessage;
 import org.apache.flink.util.SerializedThrowable;
 
 import io.fabric8.kubernetes.client.KubernetesClient;
@@ -126,6 +130,188 @@ public class JobStatusObserverTest extends 
OperatorTestBase {
                         .getState());
     }
 
+    @Test
+    void testSuspendCompletedWhenCancellationObserved() throws Exception {
+        var sessionJob = initSessionJob();
+        var status = sessionJob.getStatus();
+        var jobStatus = status.getJobStatus();
+        var jobId = JobID.fromHexString(jobStatus.getJobId());
+        FlinkResourceContext<AbstractFlinkResource<?, ?>> ctx =
+                getResourceContext(
+                        sessionJob,
+                        
TestUtils.createContextWithReadyFlinkDeployment(kubernetesClient));
+        flinkService.submitJobToSessionCluster(
+                sessionJob.getMetadata(),
+                sessionJob.getSpec(),
+                jobId,
+                ctx.getDeployConfig(sessionJob.getSpec()),
+                null);
+        flinkService.cancelJob(jobId, false);
+
+        simulateAsyncSuspendInProgress(sessionJob);
+
+        observer.observe(ctx);
+
+        assertEquals(JobStatus.CANCELED, jobStatus.getState());
+        assertEquals(
+                JobState.SUSPENDED,
+                status.getReconciliationStatus()
+                        .deserializeLastReconciledSpec()
+                        .getJob()
+                        .getState());
+        // Nothing is left to deploy for the suspend, so the upgrade must be 
completed here
+        assertEquals(ReconciliationState.DEPLOYED, 
status.getReconciliationStatus().getState());
+        
assertTrue(status.getReconciliationStatus().isLastReconciledSpecStable());
+        assertEquals(ResourceLifecycleState.SUSPENDED, 
status.getLifecycleState());
+
+        assertEquals(1, flinkResourceEventCollector.events.size());
+        var event = flinkResourceEventCollector.events.poll();
+        assertEquals(EventRecorder.Reason.JobStatusChanged.name(), 
event.getReason());
+        assertEquals("Job status changed from CANCELLING to CANCELED", 
event.getMessage());
+    }
+
+    @Test
+    void 
testSuspendedStateSurvivesJobManagerExceptionObservationAfterFinalize() throws 
Exception {
+        // After finalizeSuspendedUpgrade() completes the suspend, observe() 
may still call
+        // observeJobManagerExceptions() in the same pass (it checks the 
previous job status,
+        // not the reconciliation state). Verify this does not clobber the 
completed suspend.
+        var sessionJob = initSessionJob();
+        var status = sessionJob.getStatus();
+        var jobStatus = status.getJobStatus();
+        var jobId = JobID.fromHexString(jobStatus.getJobId());
+        FlinkResourceContext<AbstractFlinkResource<?, ?>> ctx =
+                getResourceContext(
+                        sessionJob,
+                        
TestUtils.createContextWithReadyFlinkDeployment(kubernetesClient));
+        flinkService.submitJobToSessionCluster(
+                sessionJob.getMetadata(),
+                sessionJob.getSpec(),
+                jobId,
+                ctx.getDeployConfig(sessionJob.getSpec()),
+                null);
+        flinkService.cancelJob(jobId, false);
+
+        // Pre-seed the exception cache and history so 
observeJobManagerExceptions()
+        // runs its full body instead of returning early on a null history.
+        ctx.getExceptionCacheEntry().setInitialized(true);
+        ctx.getExceptionCacheEntry().setJobId(jobId.toHexString());
+        
ctx.getExceptionCacheEntry().setLastTimestamp(Instant.ofEpochMilli(500L));
+        flinkService.addExceptionHistory(jobId, "SomeException", "trace", 
1000L);
+
+        simulateAsyncSuspendInProgress(sessionJob);
+
+        observer.observe(ctx);
+
+        // Check observeJobManagerExceptions() really executed.
+        var events =
+                kubernetesClient
+                        .v1()
+                        .events()
+                        .inNamespace(sessionJob.getMetadata().getNamespace())
+                        .list()
+                        .getItems();
+        assertTrue(
+                events.stream()
+                        .anyMatch(
+                                e ->
+                                        
EventRecorder.Reason.JobException.name()
+                                                .equals(e.getReason())));
+        assertEquals(Instant.ofEpochMilli(1000L), 
ctx.getExceptionCacheEntry().getLastTimestamp());
+
+        // The suspend completion performed by finalizeSuspendedUpgrade() must 
not be undone or
+        // altered by the subsequent observeJobManagerExceptions() call.
+        assertEquals(JobStatus.CANCELED, jobStatus.getState());
+        assertEquals(
+                JobState.SUSPENDED,
+                status.getReconciliationStatus()
+                        .deserializeLastReconciledSpec()
+                        .getJob()
+                        .getState());
+        assertEquals(ReconciliationState.DEPLOYED, 
status.getReconciliationStatus().getState());
+        
assertTrue(status.getReconciliationStatus().isLastReconciledSpecStable());
+        assertEquals(ResourceLifecycleState.SUSPENDED, 
status.getLifecycleState());
+        // The stable spec is the same document as the reconciled spec (that's 
what "stable"
+        // means here) and it too reflects the completed suspend, not some 
clobbered value.
+        assertEquals(
+                status.getReconciliationStatus().getLastReconciledSpec(),
+                status.getReconciliationStatus().getLastStableSpec());
+        assertEquals(
+                JobState.SUSPENDED,
+                
status.getReconciliationStatus().deserializeLastStableSpec().getJob().getState());
+    }
+
+    @Test
+    void testUpgradeNotCompletedWhenTargetStateStillRunning() throws Exception 
{
+        var deployment = initDeployment();
+        var status = deployment.getStatus();
+        var jobStatus = status.getJobStatus();
+        FlinkResourceContext<AbstractFlinkResource<?, ?>> ctx = 
getResourceContext(deployment);
+        flinkService.submitApplicationCluster(
+                deployment.getSpec().getJob(), 
ctx.getDeployConfig(deployment.getSpec()), false);
+        flinkService.cancelJob(JobID.fromHexString(jobStatus.getJobId()), 
false);
+
+        // The job was cancelled as the first step of an upgrade, the target 
state stays RUNNING
+        assertEquals(JobState.RUNNING, 
deployment.getSpec().getJob().getState());
+        
status.getReconciliationStatus().setState(ReconciliationState.UPGRADING);
+        jobStatus.setState(JobStatus.CANCELLING);
+
+        observer.observe(ctx);
+
+        assertEquals(
+                JobState.SUSPENDED,
+                status.getReconciliationStatus()
+                        .deserializeLastReconciledSpec()
+                        .getJob()
+                        .getState());
+        // The new spec still has to be deployed, the upgrade must not be 
completed here
+        assertEquals(ReconciliationState.UPGRADING, 
status.getReconciliationStatus().getState());
+        
assertFalse(status.getReconciliationStatus().isLastReconciledSpecStable());
+    }
+
+    @ParameterizedTest
+    @EnumSource(
+            value = JobStatus.class,
+            mode = EnumSource.Mode.INCLUDE,
+            names = {"FINISHED", "FAILED"})
+    void testSuspendCompletedWhenJobEndsWhileCancelling(JobStatus 
terminalStatus) throws Exception {
+        var deployment = initDeployment();
+        var status = deployment.getStatus();
+        FlinkResourceContext<AbstractFlinkResource<?, ?>> ctx = 
getResourceContext(deployment);
+        flinkService.submitApplicationCluster(
+                deployment.getSpec().getJob(), 
ctx.getDeployConfig(deployment.getSpec()), false);
+
+        // The job reached a globally terminal state on its own while the 
cancellation was still in
+        // flight, which the enclosing branch also treats as a completed 
cancellation
+        var job = flinkService.listJobs().get(0);
+        flinkService
+                .listJobs()
+                .set(
+                        0,
+                        Tuple3.of(
+                                job.f0,
+                                new JobStatusMessage(
+                                        job.f1.getJobId(),
+                                        job.f1.getJobName(),
+                                        terminalStatus,
+                                        job.f1.getStartTime()),
+                                job.f2));
+
+        simulateAsyncSuspendInProgress(deployment);
+
+        observer.observe(ctx);
+
+        assertEquals(
+                JobState.SUSPENDED,
+                status.getReconciliationStatus()
+                        .deserializeLastReconciledSpec()
+                        .getJob()
+                        .getState());
+        // The suspend is over regardless of how the job ended, so it must be 
completed here too
+        assertEquals(ReconciliationState.DEPLOYED, 
status.getReconciliationStatus().getState());
+        
assertTrue(status.getReconciliationStatus().isLastReconciledSpecStable());
+        assertEquals(ResourceLifecycleState.SUSPENDED, 
status.getLifecycleState());
+    }
+
     @Test
     void testFailed() throws Exception {
         var observer = new JobStatusObserver<>(eventRecorder);
@@ -769,4 +955,12 @@ public class JobStatusObserverTest extends 
OperatorTestBase {
                 .serializeAndSetLastReconciledSpec(job.getSpec(), job);
         return job;
     }
+
+    private void simulateAsyncSuspendInProgress(AbstractFlinkResource<?, ?> 
resource) {
+        // A stateless or last-state session job suspend goes through an 
asynchronous cancellation,
+        // the reconciler only initiates it, records UPGRADING and waits for 
the observer
+        resource.getSpec().getJob().setState(JobState.SUSPENDED);
+        
resource.getStatus().getReconciliationStatus().setState(ReconciliationState.UPGRADING);
+        resource.getStatus().getJobStatus().setState(JobStatus.CANCELLING);
+    }
 }

Reply via email to