-
Notifications
You must be signed in to change notification settings - Fork 243
Fix recordHeartbeat: swallowed cancel/reset/pause exceptions #2990
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
b867534
f357900
cb0373a
7c9af2d
4881e2b
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -1,6 +1,6 @@ | ||
| package io.temporal.activity; | ||
|
|
||
| import io.temporal.failure.CanceledFailure; | ||
| import io.temporal.client.ActivityCompletionException; | ||
| import javax.annotation.Nonnull; | ||
| import javax.annotation.Nullable; | ||
|
|
||
|
|
@@ -30,8 +30,14 @@ public interface ManualActivityCompletionClient { | |
| * Records heartbeat for an activity | ||
| * | ||
| * @param details to record with the heartbeat | ||
| * @throws ActivityCompletionException if the server reports the activity was cancelled, reset, or | ||
| * paused ({@link io.temporal.client.ActivityCanceledException}, {@link | ||
| * io.temporal.client.ActivityResetException}, {@link | ||
| * io.temporal.client.ActivityPausedException}), or if the heartbeat RPC fails ({@link | ||
| * io.temporal.client.ActivityCompletionFailureException}, {@link | ||
| * io.temporal.client.ActivityNotExistsException}). | ||
| */ | ||
| void recordHeartbeat(@Nullable Object details) throws CanceledFailure; | ||
| void recordHeartbeat(@Nullable Object details) throws ActivityCompletionException; | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I'll note both the old code and this are not really best practice since both are
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I could get rid of it, but we do the same thing in |
||
|
|
||
| /** | ||
| * Confirms successful cancellation to the server. | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,33 @@ | ||
| package io.temporal.internal.client; | ||
|
|
||
| import io.temporal.api.workflowservice.v1.RecordActivityTaskHeartbeatByIdResponse; | ||
| import io.temporal.api.workflowservice.v1.RecordActivityTaskHeartbeatResponse; | ||
|
|
||
| /** | ||
| * Container class to deduplicate {@link RecordActivityTaskHeartbeatByIdResponse} and {@link | ||
| * RecordActivityTaskHeartbeatResponse}. | ||
| */ | ||
| public final class ActivityHeartbeatResponse { | ||
| private final boolean cancelRequested; | ||
| private final boolean activityReset; | ||
| private final boolean activityPaused; | ||
|
|
||
| ActivityHeartbeatResponse( | ||
| boolean cancelRequested, boolean activityReset, boolean activityPaused) { | ||
| this.cancelRequested = cancelRequested; | ||
| this.activityReset = activityReset; | ||
| this.activityPaused = activityPaused; | ||
| } | ||
|
|
||
| public boolean getCancelRequested() { | ||
| return cancelRequested; | ||
| } | ||
|
|
||
| public boolean getActivityReset() { | ||
| return activityReset; | ||
| } | ||
|
|
||
| public boolean getActivityPaused() { | ||
| return activityPaused; | ||
| } | ||
| } |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,157 @@ | ||
| package io.temporal.internal.client.external; | ||
|
|
||
| import static org.junit.Assert.assertEquals; | ||
| import static org.junit.Assert.assertThrows; | ||
| import static org.junit.Assert.assertTrue; | ||
| import static org.mockito.ArgumentMatchers.any; | ||
| import static org.mockito.Mockito.mock; | ||
| import static org.mockito.Mockito.when; | ||
|
|
||
| import com.uber.m3.tally.NoopScope; | ||
| import io.grpc.Status; | ||
| import io.grpc.StatusRuntimeException; | ||
| import io.temporal.api.common.v1.WorkflowExecution; | ||
| import io.temporal.api.workflowservice.v1.RecordActivityTaskHeartbeatByIdResponse; | ||
| import io.temporal.api.workflowservice.v1.RecordActivityTaskHeartbeatResponse; | ||
| import io.temporal.api.workflowservice.v1.WorkflowServiceGrpc; | ||
| import io.temporal.client.ActivityCanceledException; | ||
| import io.temporal.client.ActivityCompletionFailureException; | ||
| import io.temporal.client.ActivityNotExistsException; | ||
| import io.temporal.client.ActivityPausedException; | ||
| import io.temporal.client.ActivityResetException; | ||
| import io.temporal.common.converter.GlobalDataConverter; | ||
| import io.temporal.serviceclient.WorkflowServiceStubs; | ||
| import io.temporal.serviceclient.WorkflowServiceStubsOptions; | ||
| import org.junit.Before; | ||
| import org.junit.Test; | ||
|
|
||
| public class ManualActivityCompletionClientImplTest { | ||
|
|
||
| private WorkflowServiceStubs service; | ||
| private WorkflowServiceGrpc.WorkflowServiceBlockingStub blockingStub; | ||
|
|
||
| @Before | ||
| public void setUp() { | ||
| service = mock(WorkflowServiceStubs.class); | ||
| blockingStub = mock(WorkflowServiceGrpc.WorkflowServiceBlockingStub.class); | ||
| when(service.blockingStub()).thenReturn(blockingStub); | ||
| when(blockingStub.withOption(any(), any())).thenReturn(blockingStub); | ||
| when(service.getServerCapabilities()) | ||
| .thenReturn( | ||
| () -> | ||
| io.temporal.api.workflowservice.v1.GetSystemInfoResponse.Capabilities | ||
| .getDefaultInstance()); | ||
| when(service.getOptions()) | ||
| .thenReturn(WorkflowServiceStubsOptions.newBuilder().validateAndBuildWithDefaults()); | ||
| } | ||
|
|
||
| private ManualActivityCompletionClientImpl clientWithTaskToken() { | ||
| return new ManualActivityCompletionClientImpl( | ||
| service, | ||
| "test-namespace", | ||
| "test-identity", | ||
| GlobalDataConverter.get(), | ||
| new NoopScope(), | ||
| new byte[] {1, 2, 3}, | ||
| null, | ||
| null, | ||
| null); | ||
| } | ||
|
|
||
| private ManualActivityCompletionClientImpl clientWithActivityId() { | ||
| return new ManualActivityCompletionClientImpl( | ||
| service, | ||
| "test-namespace", | ||
| "test-identity", | ||
| GlobalDataConverter.get(), | ||
| new NoopScope(), | ||
| null, | ||
| WorkflowExecution.newBuilder().setWorkflowId("wf").setRunId("run").build(), | ||
| "test-activity-id", | ||
| null); | ||
| } | ||
|
|
||
| @Test | ||
| public void cancelRequestedThrowsActivityCanceledExceptionNotSwallowed() { | ||
| when(blockingStub.recordActivityTaskHeartbeat(any())) | ||
| .thenReturn( | ||
| RecordActivityTaskHeartbeatResponse.newBuilder().setCancelRequested(true).build()); | ||
|
|
||
| assertThrows( | ||
| ActivityCanceledException.class, () -> clientWithTaskToken().recordHeartbeat("details")); | ||
| } | ||
|
|
||
| @Test | ||
| public void activityResetThrowsActivityResetExceptionNotSwallowed() { | ||
| when(blockingStub.recordActivityTaskHeartbeat(any())) | ||
| .thenReturn( | ||
| RecordActivityTaskHeartbeatResponse.newBuilder().setActivityReset(true).build()); | ||
|
|
||
| assertThrows( | ||
| ActivityResetException.class, () -> clientWithTaskToken().recordHeartbeat("details")); | ||
| } | ||
|
|
||
| @Test | ||
| public void activityPausedThrowsActivityPausedExceptionNotSwallowed() { | ||
| when(blockingStub.recordActivityTaskHeartbeat(any())) | ||
| .thenReturn( | ||
| RecordActivityTaskHeartbeatResponse.newBuilder().setActivityPaused(true).build()); | ||
|
|
||
| assertThrows( | ||
| ActivityPausedException.class, () -> clientWithTaskToken().recordHeartbeat("details")); | ||
| } | ||
|
|
||
| @Test | ||
| public void byIdCancelRequestedThrowsActivityCanceledExceptionNotSwallowed() { | ||
| when(blockingStub.recordActivityTaskHeartbeatById(any())) | ||
| .thenReturn( | ||
| RecordActivityTaskHeartbeatByIdResponse.newBuilder().setCancelRequested(true).build()); | ||
|
|
||
| assertThrows( | ||
| ActivityCanceledException.class, () -> clientWithActivityId().recordHeartbeat("details")); | ||
| } | ||
|
|
||
| @Test | ||
| public void byIdActivityResetThrowsActivityResetExceptionNotSwallowed() { | ||
| when(blockingStub.recordActivityTaskHeartbeatById(any())) | ||
| .thenReturn( | ||
| RecordActivityTaskHeartbeatByIdResponse.newBuilder().setActivityReset(true).build()); | ||
|
|
||
| assertThrows( | ||
| ActivityResetException.class, () -> clientWithActivityId().recordHeartbeat("details")); | ||
| } | ||
|
|
||
| @Test | ||
| public void byIdActivityPausedThrowsActivityPausedExceptionNotSwallowed() { | ||
| when(blockingStub.recordActivityTaskHeartbeatById(any())) | ||
| .thenReturn( | ||
| RecordActivityTaskHeartbeatByIdResponse.newBuilder().setActivityPaused(true).build()); | ||
|
|
||
| assertThrows( | ||
| ActivityPausedException.class, () -> clientWithActivityId().recordHeartbeat("details")); | ||
| } | ||
|
|
||
| @Test | ||
| public void notFoundIsReportedAsActivityNotExistsException() { | ||
| when(blockingStub.recordActivityTaskHeartbeat(any())) | ||
| .thenThrow(new StatusRuntimeException(Status.NOT_FOUND)); | ||
|
|
||
| assertThrows( | ||
| ActivityNotExistsException.class, () -> clientWithTaskToken().recordHeartbeat("details")); | ||
| } | ||
|
|
||
| @Test | ||
| public void rpcErrorIsReportedAsActivityCompletionFailureException() { | ||
| when(blockingStub.recordActivityTaskHeartbeat(any())) | ||
| .thenThrow(new StatusRuntimeException(Status.INTERNAL)); | ||
|
|
||
| ActivityCompletionFailureException failure = | ||
| assertThrows( | ||
| ActivityCompletionFailureException.class, | ||
| () -> clientWithTaskToken().recordHeartbeat("details")); | ||
|
|
||
| assertTrue(failure.getCause() instanceof StatusRuntimeException); | ||
| assertEquals( | ||
| Status.Code.INTERNAL, ((StatusRuntimeException) failure.getCause()).getStatus().getCode()); | ||
| } | ||
| } |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
nit: This is the docs we use elsewhere. It is easier to maintain and then move the details here into
ActivityCompletionExceptionsince there are a lot of methods that would need all these detauilsThere was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Hmm, I'll just delete the unnecessary doc and base it on
ActivityCompletionClient. I don't think there's much point in havingActivityCompletionExceptiondocument its own subclasses.