diff --git a/src/main/java/com/uber/cadence/internal/compatibility/proto/mappers/EnumMapper.java b/src/main/java/com/uber/cadence/internal/compatibility/proto/mappers/EnumMapper.java index c2179b445..e917becc8 100644 --- a/src/main/java/com/uber/cadence/internal/compatibility/proto/mappers/EnumMapper.java +++ b/src/main/java/com/uber/cadence/internal/compatibility/proto/mappers/EnumMapper.java @@ -613,4 +613,80 @@ public static com.uber.cadence.CronOverlapPolicy cronOverlapPolicy(CronOverlapPo } throw new IllegalArgumentException("unexpected enum value"); } + + public static ScheduleOverlapPolicy scheduleOverlapPolicy( + com.uber.cadence.ScheduleOverlapPolicy t) { + if (t == null) { + return ScheduleOverlapPolicy.SCHEDULE_OVERLAP_POLICY_INVALID; + } + switch (t) { + case INVALID: + return ScheduleOverlapPolicy.SCHEDULE_OVERLAP_POLICY_INVALID; + case SKIP_NEW: + return ScheduleOverlapPolicy.SCHEDULE_OVERLAP_POLICY_SKIP_NEW; + case BUFFER: + return ScheduleOverlapPolicy.SCHEDULE_OVERLAP_POLICY_BUFFER; + case CONCURRENT: + return ScheduleOverlapPolicy.SCHEDULE_OVERLAP_POLICY_CONCURRENT; + case CANCEL_PREVIOUS: + return ScheduleOverlapPolicy.SCHEDULE_OVERLAP_POLICY_CANCEL_PREVIOUS; + case TERMINATE_PREVIOUS: + return ScheduleOverlapPolicy.SCHEDULE_OVERLAP_POLICY_TERMINATE_PREVIOUS; + } + throw new IllegalArgumentException("unexpected enum value"); + } + + public static com.uber.cadence.ScheduleOverlapPolicy scheduleOverlapPolicy( + ScheduleOverlapPolicy t) { + switch (t) { + case SCHEDULE_OVERLAP_POLICY_INVALID: + case UNRECOGNIZED: + return com.uber.cadence.ScheduleOverlapPolicy.INVALID; + case SCHEDULE_OVERLAP_POLICY_SKIP_NEW: + return com.uber.cadence.ScheduleOverlapPolicy.SKIP_NEW; + case SCHEDULE_OVERLAP_POLICY_BUFFER: + return com.uber.cadence.ScheduleOverlapPolicy.BUFFER; + case SCHEDULE_OVERLAP_POLICY_CONCURRENT: + return com.uber.cadence.ScheduleOverlapPolicy.CONCURRENT; + case SCHEDULE_OVERLAP_POLICY_CANCEL_PREVIOUS: + return com.uber.cadence.ScheduleOverlapPolicy.CANCEL_PREVIOUS; + case SCHEDULE_OVERLAP_POLICY_TERMINATE_PREVIOUS: + return com.uber.cadence.ScheduleOverlapPolicy.TERMINATE_PREVIOUS; + } + throw new IllegalArgumentException("unexpected enum value"); + } + + public static ScheduleCatchUpPolicy scheduleCatchUpPolicy( + com.uber.cadence.ScheduleCatchUpPolicy t) { + if (t == null) { + return ScheduleCatchUpPolicy.SCHEDULE_CATCH_UP_POLICY_INVALID; + } + switch (t) { + case INVALID: + return ScheduleCatchUpPolicy.SCHEDULE_CATCH_UP_POLICY_INVALID; + case SKIP: + return ScheduleCatchUpPolicy.SCHEDULE_CATCH_UP_POLICY_SKIP; + case ONE: + return ScheduleCatchUpPolicy.SCHEDULE_CATCH_UP_POLICY_ONE; + case ALL: + return ScheduleCatchUpPolicy.SCHEDULE_CATCH_UP_POLICY_ALL; + } + throw new IllegalArgumentException("unexpected enum value"); + } + + public static com.uber.cadence.ScheduleCatchUpPolicy scheduleCatchUpPolicy( + ScheduleCatchUpPolicy t) { + switch (t) { + case SCHEDULE_CATCH_UP_POLICY_INVALID: + case UNRECOGNIZED: + return com.uber.cadence.ScheduleCatchUpPolicy.INVALID; + case SCHEDULE_CATCH_UP_POLICY_SKIP: + return com.uber.cadence.ScheduleCatchUpPolicy.SKIP; + case SCHEDULE_CATCH_UP_POLICY_ONE: + return com.uber.cadence.ScheduleCatchUpPolicy.ONE; + case SCHEDULE_CATCH_UP_POLICY_ALL: + return com.uber.cadence.ScheduleCatchUpPolicy.ALL; + } + throw new IllegalArgumentException("unexpected enum value"); + } } diff --git a/src/main/java/com/uber/cadence/internal/compatibility/proto/mappers/RequestMapper.java b/src/main/java/com/uber/cadence/internal/compatibility/proto/mappers/RequestMapper.java index 987d20653..1705fe1c6 100644 --- a/src/main/java/com/uber/cadence/internal/compatibility/proto/mappers/RequestMapper.java +++ b/src/main/java/com/uber/cadence/internal/compatibility/proto/mappers/RequestMapper.java @@ -23,6 +23,8 @@ import static com.uber.cadence.internal.compatibility.proto.mappers.EnumMapper.queryConsistencyLevel; import static com.uber.cadence.internal.compatibility.proto.mappers.EnumMapper.queryRejectCondition; import static com.uber.cadence.internal.compatibility.proto.mappers.EnumMapper.queryTaskCompletedType; +import static com.uber.cadence.internal.compatibility.proto.mappers.EnumMapper.scheduleCatchUpPolicy; +import static com.uber.cadence.internal.compatibility.proto.mappers.EnumMapper.scheduleOverlapPolicy; import static com.uber.cadence.internal.compatibility.proto.mappers.EnumMapper.taskListType; import static com.uber.cadence.internal.compatibility.proto.mappers.EnumMapper.workflowIdReusePolicy; import static com.uber.cadence.internal.compatibility.proto.mappers.Helpers.arrayToByteString; @@ -30,6 +32,7 @@ import static com.uber.cadence.internal.compatibility.proto.mappers.Helpers.newFieldMask; import static com.uber.cadence.internal.compatibility.proto.mappers.Helpers.nullToEmpty; import static com.uber.cadence.internal.compatibility.proto.mappers.Helpers.secondsToDuration; +import static com.uber.cadence.internal.compatibility.proto.mappers.Helpers.unixNanoToTime; import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.activeClusterSelectionPolicy; import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.activeClusters; import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.badBinaries; @@ -39,6 +42,9 @@ import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.memo; import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.payload; import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.retryPolicy; +import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.scheduleAction; +import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.schedulePolicies; +import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.scheduleSpec; import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.searchAttributes; import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.startTimeFilter; import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.statusFilter; @@ -1011,4 +1017,123 @@ public static RefreshWorkflowTasksRequest refreshWorkflowTasksRequest( .setDomain(request.getDomain() != null ? request.getDomain() : "") .build(); } + + public static com.uber.cadence.api.v1.CreateScheduleRequest createScheduleRequest( + com.uber.cadence.CreateScheduleRequest t) { + if (t == null) { + return null; + } + com.uber.cadence.api.v1.CreateScheduleRequest.Builder b = + com.uber.cadence.api.v1.CreateScheduleRequest.newBuilder() + .setDomain(nullToEmpty(t.getDomain())) + .setScheduleId(nullToEmpty(t.getScheduleId())); + if (t.getSpec() != null) b.setSpec(scheduleSpec(t.getSpec())); + if (t.getAction() != null) b.setAction(scheduleAction(t.getAction())); + if (t.getPolicies() != null) b.setPolicies(schedulePolicies(t.getPolicies())); + if (t.getMemo() != null) b.setMemo(memo(t.getMemo())); + if (t.getSearchAttributes() != null) + b.setSearchAttributes(searchAttributes(t.getSearchAttributes())); + return b.build(); + } + + public static com.uber.cadence.api.v1.DescribeScheduleRequest describeScheduleRequest( + com.uber.cadence.DescribeScheduleRequest t) { + if (t == null) { + return null; + } + return com.uber.cadence.api.v1.DescribeScheduleRequest.newBuilder() + .setDomain(nullToEmpty(t.getDomain())) + .setScheduleId(nullToEmpty(t.getScheduleId())) + .build(); + } + + public static com.uber.cadence.api.v1.UpdateScheduleRequest updateScheduleRequest( + com.uber.cadence.UpdateScheduleRequest t) { + if (t == null) { + return null; + } + com.uber.cadence.api.v1.UpdateScheduleRequest.Builder b = + com.uber.cadence.api.v1.UpdateScheduleRequest.newBuilder() + .setDomain(nullToEmpty(t.getDomain())) + .setScheduleId(nullToEmpty(t.getScheduleId())); + if (t.getSpec() != null) b.setSpec(scheduleSpec(t.getSpec())); + if (t.getAction() != null) b.setAction(scheduleAction(t.getAction())); + if (t.getPolicies() != null) b.setPolicies(schedulePolicies(t.getPolicies())); + if (t.getSearchAttributes() != null) + b.setSearchAttributes(searchAttributes(t.getSearchAttributes())); + return b.build(); + } + + public static com.uber.cadence.api.v1.DeleteScheduleRequest deleteScheduleRequest( + com.uber.cadence.DeleteScheduleRequest t) { + if (t == null) { + return null; + } + return com.uber.cadence.api.v1.DeleteScheduleRequest.newBuilder() + .setDomain(nullToEmpty(t.getDomain())) + .setScheduleId(nullToEmpty(t.getScheduleId())) + .build(); + } + + public static com.uber.cadence.api.v1.PauseScheduleRequest pauseScheduleRequest( + com.uber.cadence.PauseScheduleRequest t) { + if (t == null) { + return null; + } + return com.uber.cadence.api.v1.PauseScheduleRequest.newBuilder() + .setDomain(nullToEmpty(t.getDomain())) + .setScheduleId(nullToEmpty(t.getScheduleId())) + .setReason(nullToEmpty(t.getReason())) + .setIdentity(nullToEmpty(t.getIdentity())) + .build(); + } + + public static com.uber.cadence.api.v1.UnpauseScheduleRequest unpauseScheduleRequest( + com.uber.cadence.UnpauseScheduleRequest t) { + if (t == null) { + return null; + } + com.uber.cadence.api.v1.UnpauseScheduleRequest.Builder b = + com.uber.cadence.api.v1.UnpauseScheduleRequest.newBuilder() + .setDomain(nullToEmpty(t.getDomain())) + .setScheduleId(nullToEmpty(t.getScheduleId())) + .setReason(nullToEmpty(t.getReason())); + if (t.getCatchUpPolicy() != null) { + b.setCatchUpPolicy(scheduleCatchUpPolicy(t.getCatchUpPolicy())); + } + return b.build(); + } + + public static com.uber.cadence.api.v1.BackfillScheduleRequest backfillScheduleRequest( + com.uber.cadence.BackfillScheduleRequest t) { + if (t == null) { + return null; + } + com.uber.cadence.api.v1.BackfillScheduleRequest.Builder b = + com.uber.cadence.api.v1.BackfillScheduleRequest.newBuilder() + .setDomain(nullToEmpty(t.getDomain())) + .setScheduleId(nullToEmpty(t.getScheduleId())) + .setStartTime(unixNanoToTime(t.getStartTimeNano())) + .setEndTime(unixNanoToTime(t.getEndTimeNano())); + if (t.getOverlapPolicy() != null) { + b.setOverlapPolicy(scheduleOverlapPolicy(t.getOverlapPolicy())); + } + if (t.getBackfillId() != null) { + b.setBackfillId(t.getBackfillId()); + } + return b.build(); + } + + public static com.uber.cadence.api.v1.ListSchedulesRequest listSchedulesRequest( + com.uber.cadence.ListSchedulesRequest t) { + if (t == null) { + return null; + } + com.uber.cadence.api.v1.ListSchedulesRequest.Builder b = + com.uber.cadence.api.v1.ListSchedulesRequest.newBuilder() + .setDomain(nullToEmpty(t.getDomain())) + .setPageSize(t.getPageSize()); + if (t.getNextPageToken() != null) b.setNextPageToken(arrayToByteString(t.getNextPageToken())); + return b.build(); + } } diff --git a/src/main/java/com/uber/cadence/internal/compatibility/proto/mappers/ResponseMapper.java b/src/main/java/com/uber/cadence/internal/compatibility/proto/mappers/ResponseMapper.java index 71fbf67cd..fbbf257ee 100644 --- a/src/main/java/com/uber/cadence/internal/compatibility/proto/mappers/ResponseMapper.java +++ b/src/main/java/com/uber/cadence/internal/compatibility/proto/mappers/ResponseMapper.java @@ -30,12 +30,20 @@ import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.describeDomainResponseArray; import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.header; import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.indexedValueTypeMap; +import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.memo; import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.payload; import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.pendingActivityInfoArray; import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.pendingChildExecutionInfoArray; import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.pendingDecisionInfo; import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.pollerInfoArray; import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.queryRejected; +import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.scheduleAction; +import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.scheduleInfo; +import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.scheduleListEntry; +import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.schedulePolicies; +import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.scheduleSpec; +import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.scheduleState; +import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.searchAttributes; import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.supportedClientVersions; import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.taskList; import static com.uber.cadence.internal.compatibility.proto.mappers.TypeMapper.taskListPartitionMetadataArray; @@ -518,4 +526,89 @@ public static com.uber.cadence.GetTaskListsByDomainResponse getTaskListsByDomain Collectors.toMap(Map.Entry::getKey, e -> describeTaskListResponse(e.getValue())))); return res; } + + public static com.uber.cadence.CreateScheduleResponse createScheduleResponse( + com.uber.cadence.api.v1.CreateScheduleResponse t) { + if (t == null) { + return null; + } + com.uber.cadence.CreateScheduleResponse res = new com.uber.cadence.CreateScheduleResponse(); + res.setScheduleId(t.getScheduleId()); + return res; + } + + public static com.uber.cadence.DescribeScheduleResponse describeScheduleResponse( + com.uber.cadence.api.v1.DescribeScheduleResponse t) { + if (t == null) { + return null; + } + com.uber.cadence.DescribeScheduleResponse res = new com.uber.cadence.DescribeScheduleResponse(); + res.setSpec(scheduleSpec(t.getSpec())); + res.setAction(scheduleAction(t.getAction())); + res.setPolicies(schedulePolicies(t.getPolicies())); + res.setState(scheduleState(t.getState())); + res.setInfo(scheduleInfo(t.getInfo())); + res.setMemo(memo(t.getMemo())); + res.setSearchAttributes(searchAttributes(t.getSearchAttributes())); + return res; + } + + public static com.uber.cadence.ListSchedulesResponse listSchedulesResponse( + com.uber.cadence.api.v1.ListSchedulesResponse t) { + if (t == null) { + return null; + } + com.uber.cadence.ListSchedulesResponse res = new com.uber.cadence.ListSchedulesResponse(); + if (t.getSchedulesCount() > 0) { + java.util.List entries = new java.util.ArrayList<>(); + for (com.uber.cadence.api.v1.ScheduleListEntry e : t.getSchedulesList()) { + entries.add(scheduleListEntry(e)); + } + res.setSchedules(entries); + } + if (t.getNextPageToken() != null && !t.getNextPageToken().isEmpty()) { + res.setNextPageToken(byteStringToArray(t.getNextPageToken())); + } + return res; + } + + public static com.uber.cadence.UpdateScheduleResponse updateScheduleResponse( + com.uber.cadence.api.v1.UpdateScheduleResponse t) { + if (t == null) { + return null; + } + return new com.uber.cadence.UpdateScheduleResponse(); + } + + public static com.uber.cadence.DeleteScheduleResponse deleteScheduleResponse( + com.uber.cadence.api.v1.DeleteScheduleResponse t) { + if (t == null) { + return null; + } + return new com.uber.cadence.DeleteScheduleResponse(); + } + + public static com.uber.cadence.PauseScheduleResponse pauseScheduleResponse( + com.uber.cadence.api.v1.PauseScheduleResponse t) { + if (t == null) { + return null; + } + return new com.uber.cadence.PauseScheduleResponse(); + } + + public static com.uber.cadence.UnpauseScheduleResponse unpauseScheduleResponse( + com.uber.cadence.api.v1.UnpauseScheduleResponse t) { + if (t == null) { + return null; + } + return new com.uber.cadence.UnpauseScheduleResponse(); + } + + public static com.uber.cadence.BackfillScheduleResponse backfillScheduleResponse( + com.uber.cadence.api.v1.BackfillScheduleResponse t) { + if (t == null) { + return null; + } + return new com.uber.cadence.BackfillScheduleResponse(); + } } diff --git a/src/main/java/com/uber/cadence/internal/compatibility/proto/mappers/TypeMapper.java b/src/main/java/com/uber/cadence/internal/compatibility/proto/mappers/TypeMapper.java index d993f61e6..d2c4a7f76 100644 --- a/src/main/java/com/uber/cadence/internal/compatibility/proto/mappers/TypeMapper.java +++ b/src/main/java/com/uber/cadence/internal/compatibility/proto/mappers/TypeMapper.java @@ -23,6 +23,8 @@ import static com.uber.cadence.internal.compatibility.proto.mappers.EnumMapper.pendingActivityState; import static com.uber.cadence.internal.compatibility.proto.mappers.EnumMapper.pendingDecisionState; import static com.uber.cadence.internal.compatibility.proto.mappers.EnumMapper.queryResultType; +import static com.uber.cadence.internal.compatibility.proto.mappers.EnumMapper.scheduleCatchUpPolicy; +import static com.uber.cadence.internal.compatibility.proto.mappers.EnumMapper.scheduleOverlapPolicy; import static com.uber.cadence.internal.compatibility.proto.mappers.EnumMapper.taskListKind; import static com.uber.cadence.internal.compatibility.proto.mappers.EnumMapper.workflowExecutionCloseStatus; import static com.uber.cadence.internal.compatibility.proto.mappers.Helpers.arrayToByteString; @@ -1088,4 +1090,167 @@ static com.uber.cadence.ClusterAttribute clusterAttribute(ClusterAttribute t) { attr.setName(t.getName()); return attr; } + + static ScheduleSpec scheduleSpec(com.uber.cadence.ScheduleSpec t) { + if (t == null) { + return ScheduleSpec.getDefaultInstance(); + } + return ScheduleSpec.newBuilder() + .setCronExpression(Helpers.nullToEmpty(t.getCronExpression())) + .setStartTime(unixNanoToTime(t.getStartTimeNano())) + .setEndTime(unixNanoToTime(t.getEndTimeNano())) + .setJitter(secondsToDuration(t.getJitterInSeconds())) + .build(); + } + + static com.uber.cadence.ScheduleSpec scheduleSpec(ScheduleSpec t) { + if (t == null || t == ScheduleSpec.getDefaultInstance()) { + return null; + } + com.uber.cadence.ScheduleSpec res = new com.uber.cadence.ScheduleSpec(); + res.setCronExpression(t.getCronExpression()); + res.setStartTimeNano(timeToUnixNano(t.getStartTime())); + res.setEndTimeNano(timeToUnixNano(t.getEndTime())); + res.setJitterInSeconds(durationToSeconds(t.getJitter())); + return res; + } + + static ScheduleAction scheduleAction(com.uber.cadence.ScheduleAction t) { + if (t == null || t.getStartWorkflow() == null) { + return ScheduleAction.getDefaultInstance(); + } + com.uber.cadence.ScheduleStartWorkflowAction sw = t.getStartWorkflow(); + ScheduleAction.StartWorkflowAction.Builder swb = + ScheduleAction.StartWorkflowAction.newBuilder(); + if (sw.getWorkflowType() != null) swb.setWorkflowType(workflowType(sw.getWorkflowType())); + if (sw.getTaskList() != null) swb.setTaskList(taskList(sw.getTaskList())); + if (sw.getInput() != null) swb.setInput(payload(sw.getInput())); + swb.setWorkflowIdPrefix(Helpers.nullToEmpty(sw.getWorkflowIdPrefix())); + swb.setExecutionStartToCloseTimeout( + secondsToDuration(sw.getExecutionStartToCloseTimeoutSeconds())); + swb.setTaskStartToCloseTimeout(secondsToDuration(sw.getTaskStartToCloseTimeoutSeconds())); + if (sw.getRetryPolicy() != null) swb.setRetryPolicy(retryPolicy(sw.getRetryPolicy())); + if (sw.getMemo() != null) swb.setMemo(memo(sw.getMemo())); + if (sw.getSearchAttributes() != null) + swb.setSearchAttributes(searchAttributes(sw.getSearchAttributes())); + return ScheduleAction.newBuilder().setStartWorkflow(swb.build()).build(); + } + + static com.uber.cadence.ScheduleAction scheduleAction(ScheduleAction t) { + if (t == null || !t.hasStartWorkflow()) { + return null; + } + ScheduleAction.StartWorkflowAction proto = t.getStartWorkflow(); + com.uber.cadence.ScheduleStartWorkflowAction sw = + new com.uber.cadence.ScheduleStartWorkflowAction(); + sw.setWorkflowType(workflowType(proto.getWorkflowType())); + sw.setTaskList(taskList(proto.getTaskList())); + sw.setInput(payload(proto.getInput())); + sw.setWorkflowIdPrefix(proto.getWorkflowIdPrefix()); + sw.setExecutionStartToCloseTimeoutSeconds( + durationToSeconds(proto.getExecutionStartToCloseTimeout())); + sw.setTaskStartToCloseTimeoutSeconds(durationToSeconds(proto.getTaskStartToCloseTimeout())); + sw.setRetryPolicy(retryPolicy(proto.getRetryPolicy())); + sw.setMemo(memo(proto.getMemo())); + sw.setSearchAttributes(searchAttributes(proto.getSearchAttributes())); + com.uber.cadence.ScheduleAction action = new com.uber.cadence.ScheduleAction(); + action.setStartWorkflow(sw); + return action; + } + + static SchedulePolicies schedulePolicies(com.uber.cadence.SchedulePolicies t) { + if (t == null) { + return SchedulePolicies.getDefaultInstance(); + } + return SchedulePolicies.newBuilder() + .setOverlapPolicy(scheduleOverlapPolicy(t.getOverlapPolicy())) + .setCatchUpPolicy(scheduleCatchUpPolicy(t.getCatchUpPolicy())) + .setCatchUpWindow(secondsToDuration(t.getCatchUpWindowInSeconds())) + .setPauseOnFailure(t.isPauseOnFailure()) + .setBufferLimit(t.getBufferLimit()) + .setConcurrencyLimit(t.getConcurrencyLimit()) + .build(); + } + + static com.uber.cadence.SchedulePolicies schedulePolicies(SchedulePolicies t) { + if (t == null || t == SchedulePolicies.getDefaultInstance()) { + return null; + } + com.uber.cadence.SchedulePolicies res = new com.uber.cadence.SchedulePolicies(); + res.setOverlapPolicy(scheduleOverlapPolicy(t.getOverlapPolicy())); + res.setCatchUpPolicy(scheduleCatchUpPolicy(t.getCatchUpPolicy())); + res.setCatchUpWindowInSeconds(durationToSeconds(t.getCatchUpWindow())); + res.setPauseOnFailure(t.getPauseOnFailure()); + res.setBufferLimit(t.getBufferLimit()); + res.setConcurrencyLimit(t.getConcurrencyLimit()); + return res; + } + + static com.uber.cadence.SchedulePauseInfo schedulePauseInfo(SchedulePauseInfo t) { + if (t == null || t == SchedulePauseInfo.getDefaultInstance()) { + return null; + } + com.uber.cadence.SchedulePauseInfo res = new com.uber.cadence.SchedulePauseInfo(); + res.setReason(t.getReason()); + res.setPausedTimeNano(timeToUnixNano(t.getPausedAt())); + res.setPausedBy(t.getPausedBy()); + return res; + } + + static com.uber.cadence.ScheduleState scheduleState(ScheduleState t) { + if (t == null || t == ScheduleState.getDefaultInstance()) { + return null; + } + com.uber.cadence.ScheduleState res = new com.uber.cadence.ScheduleState(); + res.setPaused(t.getPaused()); + res.setPauseInfo(schedulePauseInfo(t.getPauseInfo())); + return res; + } + + static com.uber.cadence.BackfillInfo backfillInfo(BackfillInfo t) { + if (t == null || t == BackfillInfo.getDefaultInstance()) { + return null; + } + com.uber.cadence.BackfillInfo res = new com.uber.cadence.BackfillInfo(); + res.setBackfillId(t.getBackfillId()); + res.setStartTimeNano(timeToUnixNano(t.getStartTime())); + res.setEndTimeNano(timeToUnixNano(t.getEndTime())); + res.setRunsCompleted(t.getRunsCompleted()); + res.setRunsTotal(t.getRunsTotal()); + return res; + } + + static com.uber.cadence.ScheduleInfo scheduleInfo(ScheduleInfo t) { + if (t == null || t == ScheduleInfo.getDefaultInstance()) { + return null; + } + com.uber.cadence.ScheduleInfo res = new com.uber.cadence.ScheduleInfo(); + res.setLastRunTimeNano(timeToUnixNano(t.getLastRunTime())); + res.setNextRunTimeNano(timeToUnixNano(t.getNextRunTime())); + res.setTotalRuns(t.getTotalRuns()); + res.setCreateTimeNano(timeToUnixNano(t.getCreateTime())); + res.setLastUpdateTimeNano(timeToUnixNano(t.getLastUpdateTime())); + if (t.getOngoingBackfillsCount() > 0) { + List backfills = new ArrayList<>(); + for (BackfillInfo b : t.getOngoingBackfillsList()) { + backfills.add(backfillInfo(b)); + } + res.setOngoingBackfills(backfills); + } + res.setMissedRuns(t.getMissedRuns()); + res.setSkippedRuns(t.getSkippedRuns()); + return res; + } + + static com.uber.cadence.ScheduleListEntry scheduleListEntry(ScheduleListEntry t) { + if (t == null || t == ScheduleListEntry.getDefaultInstance()) { + return null; + } + com.uber.cadence.ScheduleListEntry res = new com.uber.cadence.ScheduleListEntry(); + res.setScheduleId(t.getScheduleId()); + res.setWorkflowType(workflowType(t.getWorkflowType())); + res.setState(scheduleState(t.getState())); + res.setCronExpression(t.getCronExpression()); + return res; + } } diff --git a/src/main/java/com/uber/cadence/internal/compatibility/proto/serviceclient/GrpcServiceStubs.java b/src/main/java/com/uber/cadence/internal/compatibility/proto/serviceclient/GrpcServiceStubs.java index 6e76e5abd..f0290559e 100644 --- a/src/main/java/com/uber/cadence/internal/compatibility/proto/serviceclient/GrpcServiceStubs.java +++ b/src/main/java/com/uber/cadence/internal/compatibility/proto/serviceclient/GrpcServiceStubs.java @@ -22,6 +22,7 @@ import com.uber.cadence.api.v1.MetaAPIGrpc; import com.uber.cadence.api.v1.MetaAPIGrpc.MetaAPIBlockingStub; import com.uber.cadence.api.v1.MetaAPIGrpc.MetaAPIFutureStub; +import com.uber.cadence.api.v1.ScheduleAPIGrpc; import com.uber.cadence.api.v1.VisibilityAPIGrpc; import com.uber.cadence.api.v1.VisibilityAPIGrpc.VisibilityAPIBlockingStub; import com.uber.cadence.api.v1.VisibilityAPIGrpc.VisibilityAPIFutureStub; @@ -93,6 +94,8 @@ final class GrpcServiceStubs implements IGrpcServiceStubs { private final WorkflowAPIGrpc.WorkflowAPIFutureStub workflowFutureStub; private final MetaAPIGrpc.MetaAPIBlockingStub metaBlockingStub; private final MetaAPIGrpc.MetaAPIFutureStub metaFutureStub; + private final ScheduleAPIGrpc.ScheduleAPIBlockingStub scheduleBlockingStub; + private final ScheduleAPIGrpc.ScheduleAPIFutureStub scheduleFutureStub; GrpcServiceStubs(ClientOptions options) { this.options = options; @@ -146,6 +149,8 @@ final class GrpcServiceStubs implements IGrpcServiceStubs { this.workflowFutureStub = WorkflowAPIGrpc.newFutureStub(interceptedChannel); this.metaBlockingStub = MetaAPIGrpc.newBlockingStub(interceptedChannel); this.metaFutureStub = MetaAPIGrpc.newFutureStub(interceptedChannel); + this.scheduleBlockingStub = ScheduleAPIGrpc.newBlockingStub(interceptedChannel); + this.scheduleFutureStub = ScheduleAPIGrpc.newFutureStub(interceptedChannel); } private ClientInterceptor newAuthorizationInterceptor(IAuthorizationProvider provider) { @@ -412,6 +417,16 @@ public WorkflowAPIFutureStub workflowFutureStub() { return workflowFutureStub; } + @Override + public ScheduleAPIGrpc.ScheduleAPIBlockingStub scheduleBlockingStub() { + return scheduleBlockingStub; + } + + @Override + public ScheduleAPIGrpc.ScheduleAPIFutureStub scheduleFutureStub() { + return scheduleFutureStub; + } + @Override public void shutdown() { shutdownRequested.set(true); diff --git a/src/main/java/com/uber/cadence/internal/compatibility/proto/serviceclient/IGrpcServiceStubs.java b/src/main/java/com/uber/cadence/internal/compatibility/proto/serviceclient/IGrpcServiceStubs.java index b0b0da2ec..4e4656d97 100644 --- a/src/main/java/com/uber/cadence/internal/compatibility/proto/serviceclient/IGrpcServiceStubs.java +++ b/src/main/java/com/uber/cadence/internal/compatibility/proto/serviceclient/IGrpcServiceStubs.java @@ -18,6 +18,7 @@ import com.uber.cadence.api.v1.DomainAPIGrpc; import com.uber.cadence.api.v1.MetaAPIGrpc.MetaAPIBlockingStub; import com.uber.cadence.api.v1.MetaAPIGrpc.MetaAPIFutureStub; +import com.uber.cadence.api.v1.ScheduleAPIGrpc; import com.uber.cadence.api.v1.VisibilityAPIGrpc; import com.uber.cadence.api.v1.WorkerAPIGrpc; import com.uber.cadence.api.v1.WorkflowAPIGrpc; @@ -65,6 +66,12 @@ static IGrpcServiceStubs newInstance(ClientOptions options) { /** @return Future (asynchronous) stub to workflow service. */ WorkflowAPIGrpc.WorkflowAPIFutureStub workflowFutureStub(); + /** @return Blocking (synchronous) stub to schedule service. */ + ScheduleAPIGrpc.ScheduleAPIBlockingStub scheduleBlockingStub(); + + /** @return Future (asynchronous) stub to schedule service. */ + ScheduleAPIGrpc.ScheduleAPIFutureStub scheduleFutureStub(); + /** @return Blocking (synchronous) stub to meta service. */ MetaAPIFutureStub metaFutureStub(); diff --git a/src/main/java/com/uber/cadence/internal/sync/ScheduleClientImpl.java b/src/main/java/com/uber/cadence/internal/sync/ScheduleClientImpl.java new file mode 100644 index 000000000..60afcd94c --- /dev/null +++ b/src/main/java/com/uber/cadence/internal/sync/ScheduleClientImpl.java @@ -0,0 +1,323 @@ +/* + * Copyright 2012-2016 Amazon.com, Inc. or its affiliates. All Rights Reserved. + * + * Modifications copyright (C) 2017 Uber Technologies, Inc. + * + * Licensed under the Apache License, Version 2.0 (the "License"). You may not + * use this file except in compliance with the License. A copy of the License is + * located at + * + * http://aws.amazon.com/apache2.0 + * + * or in the "license" file accompanying this file. This file 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 com.uber.cadence.internal.sync; + +import com.uber.cadence.BackfillScheduleRequest; +import com.uber.cadence.BackfillScheduleResponse; +import com.uber.cadence.CreateScheduleRequest; +import com.uber.cadence.CreateScheduleResponse; +import com.uber.cadence.DeleteScheduleRequest; +import com.uber.cadence.DeleteScheduleResponse; +import com.uber.cadence.DescribeScheduleRequest; +import com.uber.cadence.DescribeScheduleResponse; +import com.uber.cadence.ListSchedulesRequest; +import com.uber.cadence.ListSchedulesResponse; +import com.uber.cadence.PauseScheduleRequest; +import com.uber.cadence.PauseScheduleResponse; +import com.uber.cadence.UnpauseScheduleRequest; +import com.uber.cadence.UnpauseScheduleResponse; +import com.uber.cadence.UpdateScheduleRequest; +import com.uber.cadence.UpdateScheduleResponse; +import com.uber.cadence.client.ScheduleBackfill; +import com.uber.cadence.client.ScheduleClient; +import com.uber.cadence.client.schedule.ScheduleAction; +import com.uber.cadence.client.schedule.ScheduleCatchUpPolicy; +import com.uber.cadence.client.schedule.ScheduleDescription; +import com.uber.cadence.client.schedule.ScheduleInfo; +import com.uber.cadence.client.schedule.ScheduleOverlapPolicy; +import com.uber.cadence.client.schedule.SchedulePolicies; +import com.uber.cadence.client.schedule.ScheduleSpec; +import com.uber.cadence.client.schedule.ScheduleState; +import com.uber.cadence.common.RetryOptions; +import com.uber.cadence.serviceclient.IWorkflowService; +import java.time.Duration; +import java.time.Instant; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CompletableFuture; + +final class ScheduleClientImpl implements ScheduleClient { + + private final IWorkflowService service; + private final String domain; + + ScheduleClientImpl(IWorkflowService service, String domain) { + this.service = service; + this.domain = domain; + } + + @Override + public CompletableFuture createSchedule( + String scheduleId, CreateScheduleRequest request) { + request.setDomain(domain); + request.setScheduleId(scheduleId); + return service.CreateSchedule(request); + } + + @Override + public CompletableFuture describeSchedule(String scheduleId) { + DescribeScheduleRequest request = + new DescribeScheduleRequest().setDomain(domain).setScheduleId(scheduleId); + return service.DescribeSchedule(request).thenApply(ScheduleClientImpl::toScheduleDescription); + } + + @Override + public CompletableFuture updateSchedule( + String scheduleId, UpdateScheduleRequest request) { + request.setDomain(domain); + request.setScheduleId(scheduleId); + return service.UpdateSchedule(request); + } + + @Override + public CompletableFuture deleteSchedule(String scheduleId) { + DeleteScheduleRequest request = + new DeleteScheduleRequest().setDomain(domain).setScheduleId(scheduleId); + return service.DeleteSchedule(request); + } + + @Override + public CompletableFuture pauseSchedule(String scheduleId, String reason) { + PauseScheduleRequest request = + new PauseScheduleRequest().setDomain(domain).setScheduleId(scheduleId).setReason(reason); + return service.PauseSchedule(request); + } + + @Override + public CompletableFuture unpauseSchedule( + String scheduleId, String reason) { + UnpauseScheduleRequest request = + new UnpauseScheduleRequest().setDomain(domain).setScheduleId(scheduleId).setReason(reason); + return service.UnpauseSchedule(request); + } + + @Override + public CompletableFuture> backfillSchedule( + String scheduleId, List backfills) { + List> futures = new ArrayList<>(); + for (ScheduleBackfill bf : backfills) { + BackfillScheduleRequest request = + new BackfillScheduleRequest() + .setDomain(domain) + .setScheduleId(scheduleId) + .setStartTimeNano(bf.getStartTime().toEpochMilli() * 1_000_000L) + .setEndTimeNano(bf.getEndTime().toEpochMilli() * 1_000_000L); + if (bf.getOverlapPolicy() != null) { + request.setOverlapPolicy(toThriftOverlapPolicy(bf.getOverlapPolicy())); + } + futures.add(service.BackfillSchedule(request)); + } + return CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])) + .thenApply( + v -> { + List results = new ArrayList<>(); + for (CompletableFuture f : futures) { + results.add(f.join()); + } + return results; + }); + } + + @Override + public CompletableFuture listSchedules( + int pageSize, byte[] nextPageToken) { + ListSchedulesRequest request = + new ListSchedulesRequest() + .setDomain(domain) + .setPageSize(pageSize) + .setNextPageToken(nextPageToken); + return service.ListSchedules(request); + } + + private static ScheduleDescription toScheduleDescription(DescribeScheduleResponse r) { + if (r == null) { + return null; + } + return new ScheduleDescription( + toScheduleSpec(r.getSpec()), + toScheduleAction(r.getAction()), + toSchedulePolicies(r.getPolicies()), + toScheduleState(r.getState()), + toScheduleInfo(r.getInfo()), + toObjectMap(r.getMemo() != null ? r.getMemo().getFields() : null), + toObjectMap( + r.getSearchAttributes() != null ? r.getSearchAttributes().getIndexedFields() : null)); + } + + private static ScheduleSpec toScheduleSpec(com.uber.cadence.ScheduleSpec t) { + if (t == null) return null; + return ScheduleSpec.newBuilder() + .setCronExpression(t.getCronExpression()) + .setStartTime(nanosToInstant(t.getStartTimeNano())) + .setEndTime(nanosToInstant(t.getEndTimeNano())) + .setJitter(Duration.ofSeconds(t.getJitterInSeconds())) + .build(); + } + + private static ScheduleAction toScheduleAction(com.uber.cadence.ScheduleAction t) { + if (t == null || t.getStartWorkflow() == null) return null; + com.uber.cadence.ScheduleStartWorkflowAction sw = t.getStartWorkflow(); + ScheduleAction.StartWorkflowAction.Builder swb = + ScheduleAction.StartWorkflowAction.newBuilder() + .setWorkflowType(sw.getWorkflowType() != null ? sw.getWorkflowType().getName() : null) + .setTaskList(sw.getTaskList() != null ? sw.getTaskList().getName() : null) + .setInput(sw.getInput()) + .setWorkflowIdPrefix(sw.getWorkflowIdPrefix()) + .setExecutionStartToCloseTimeout( + Duration.ofSeconds(sw.getExecutionStartToCloseTimeoutSeconds())) + .setTaskStartToCloseTimeout(Duration.ofSeconds(sw.getTaskStartToCloseTimeoutSeconds())); + if (sw.getRetryPolicy() != null) { + swb.setRetryOptions(toRetryOptions(sw.getRetryPolicy())); + } + if (sw.getMemo() != null) { + swb.setMemo(toObjectMap(sw.getMemo().getFields())); + } + if (sw.getSearchAttributes() != null) { + swb.setSearchAttributes(toObjectMap(sw.getSearchAttributes().getIndexedFields())); + } + return ScheduleAction.newBuilder().setStartWorkflow(swb.build()).build(); + } + + private static SchedulePolicies toSchedulePolicies(com.uber.cadence.SchedulePolicies t) { + if (t == null) return null; + return SchedulePolicies.newBuilder() + .setOverlapPolicy(toClientOverlapPolicy(t.getOverlapPolicy())) + .setCatchUpPolicy(toClientCatchUpPolicy(t.getCatchUpPolicy())) + .setCatchUpWindow(Duration.ofSeconds(t.getCatchUpWindowInSeconds())) + .setPauseOnFailure(t.isPauseOnFailure()) + .setBufferLimit(t.getBufferLimit()) + .setConcurrencyLimit(t.getConcurrencyLimit()) + .build(); + } + + private static ScheduleState toScheduleState(com.uber.cadence.ScheduleState t) { + if (t == null) return null; + String pauseReason = null; + Instant pausedAt = null; + String pausedBy = null; + if (t.getPauseInfo() != null) { + com.uber.cadence.SchedulePauseInfo pi = t.getPauseInfo(); + pauseReason = pi.getReason(); + pausedAt = nanosToInstant(pi.getPausedTimeNano()); + pausedBy = pi.getPausedBy(); + } + return new ScheduleState(t.isPaused(), pauseReason, pausedAt, pausedBy); + } + + private static ScheduleInfo toScheduleInfo(com.uber.cadence.ScheduleInfo t) { + if (t == null) return null; + List backfills = new ArrayList<>(); + if (t.getOngoingBackfills() != null) { + for (com.uber.cadence.BackfillInfo b : t.getOngoingBackfills()) { + backfills.add( + new ScheduleInfo.BackfillInfo( + b.getBackfillId(), + nanosToInstant(b.getStartTimeNano()), + nanosToInstant(b.getEndTimeNano()), + b.getRunsCompleted(), + b.getRunsTotal())); + } + } + return new ScheduleInfo( + nanosToInstant(t.getLastRunTimeNano()), + nanosToInstant(t.getNextRunTimeNano()), + t.getTotalRuns(), + nanosToInstant(t.getCreateTimeNano()), + nanosToInstant(t.getLastUpdateTimeNano()), + backfills, + t.getMissedRuns(), + t.getSkippedRuns()); + } + + private static RetryOptions toRetryOptions(com.uber.cadence.RetryPolicy p) { + return new RetryOptions.Builder() + .setInitialInterval(Duration.ofSeconds(p.getInitialIntervalInSeconds())) + .setMaximumInterval(Duration.ofSeconds(p.getMaximumIntervalInSeconds())) + .setBackoffCoefficient(p.getBackoffCoefficient()) + .setMaximumAttempts(p.getMaximumAttempts()) + .setExpiration(Duration.ofSeconds(p.getExpirationIntervalInSeconds())) + .build(); + } + + private static ScheduleOverlapPolicy toClientOverlapPolicy( + com.uber.cadence.ScheduleOverlapPolicy t) { + if (t == null) return null; + switch (t) { + case SKIP_NEW: + return ScheduleOverlapPolicy.SKIP_NEW; + case BUFFER: + return ScheduleOverlapPolicy.BUFFER; + case CONCURRENT: + return ScheduleOverlapPolicy.CONCURRENT; + case CANCEL_PREVIOUS: + return ScheduleOverlapPolicy.CANCEL_PREVIOUS; + case TERMINATE_PREVIOUS: + return ScheduleOverlapPolicy.TERMINATE_PREVIOUS; + default: + return null; + } + } + + private static com.uber.cadence.ScheduleOverlapPolicy toThriftOverlapPolicy( + ScheduleOverlapPolicy p) { + switch (p) { + case SKIP_NEW: + return com.uber.cadence.ScheduleOverlapPolicy.SKIP_NEW; + case BUFFER: + return com.uber.cadence.ScheduleOverlapPolicy.BUFFER; + case CONCURRENT: + return com.uber.cadence.ScheduleOverlapPolicy.CONCURRENT; + case CANCEL_PREVIOUS: + return com.uber.cadence.ScheduleOverlapPolicy.CANCEL_PREVIOUS; + case TERMINATE_PREVIOUS: + return com.uber.cadence.ScheduleOverlapPolicy.TERMINATE_PREVIOUS; + default: + throw new IllegalArgumentException("unknown ScheduleOverlapPolicy: " + p); + } + } + + private static ScheduleCatchUpPolicy toClientCatchUpPolicy( + com.uber.cadence.ScheduleCatchUpPolicy t) { + if (t == null) return null; + switch (t) { + case SKIP: + return ScheduleCatchUpPolicy.SKIP; + case ONE: + return ScheduleCatchUpPolicy.ONE; + case ALL: + return ScheduleCatchUpPolicy.ALL; + default: + return null; + } + } + + private static Instant nanosToInstant(long nanos) { + if (nanos == 0) return null; + return Instant.ofEpochSecond(nanos / 1_000_000_000L, nanos % 1_000_000_000L); + } + + @SuppressWarnings("unchecked") + private static Map toObjectMap(Map src) { + if (src == null || src.isEmpty()) return null; + Map result = new HashMap<>(); + src.forEach((k, v) -> result.put(k, v)); + return result; + } +} diff --git a/src/main/java/com/uber/cadence/internal/sync/TestActivityEnvironmentInternal.java b/src/main/java/com/uber/cadence/internal/sync/TestActivityEnvironmentInternal.java index 97ad25b89..7b7f55811 100644 --- a/src/main/java/com/uber/cadence/internal/sync/TestActivityEnvironmentInternal.java +++ b/src/main/java/com/uber/cadence/internal/sync/TestActivityEnvironmentInternal.java @@ -583,48 +583,45 @@ public void RefreshWorkflowTasks(RefreshWorkflowTasksRequest request) } @Override - public CreateScheduleResponse CreateSchedule(CreateScheduleRequest request) - throws CadenceError { + public CompletableFuture CreateSchedule(CreateScheduleRequest request) { return impl.CreateSchedule(request); } @Override - public DescribeScheduleResponse DescribeSchedule(DescribeScheduleRequest request) - throws CadenceError { + public CompletableFuture DescribeSchedule( + DescribeScheduleRequest request) { return impl.DescribeSchedule(request); } @Override - public UpdateScheduleResponse UpdateSchedule(UpdateScheduleRequest request) - throws CadenceError { + public CompletableFuture UpdateSchedule(UpdateScheduleRequest request) { return impl.UpdateSchedule(request); } @Override - public DeleteScheduleResponse DeleteSchedule(DeleteScheduleRequest request) - throws CadenceError { + public CompletableFuture DeleteSchedule(DeleteScheduleRequest request) { return impl.DeleteSchedule(request); } @Override - public PauseScheduleResponse PauseSchedule(PauseScheduleRequest request) throws CadenceError { + public CompletableFuture PauseSchedule(PauseScheduleRequest request) { return impl.PauseSchedule(request); } @Override - public UnpauseScheduleResponse UnpauseSchedule(UnpauseScheduleRequest request) - throws CadenceError { + public CompletableFuture UnpauseSchedule( + UnpauseScheduleRequest request) { return impl.UnpauseSchedule(request); } @Override - public BackfillScheduleResponse BackfillSchedule(BackfillScheduleRequest request) - throws CadenceError { + public CompletableFuture BackfillSchedule( + BackfillScheduleRequest request) { return impl.BackfillSchedule(request); } @Override - public ListSchedulesResponse ListSchedules(ListSchedulesRequest request) throws CadenceError { + public CompletableFuture ListSchedules(ListSchedulesRequest request) { return impl.ListSchedules(request); } diff --git a/src/main/java/com/uber/cadence/internal/sync/TestWorkflowEnvironmentInternal.java b/src/main/java/com/uber/cadence/internal/sync/TestWorkflowEnvironmentInternal.java index 61d72ae00..2a98e585a 100644 --- a/src/main/java/com/uber/cadence/internal/sync/TestWorkflowEnvironmentInternal.java +++ b/src/main/java/com/uber/cadence/internal/sync/TestWorkflowEnvironmentInternal.java @@ -481,48 +481,45 @@ public void RefreshWorkflowTasks(RefreshWorkflowTasksRequest request) } @Override - public CreateScheduleResponse CreateSchedule(CreateScheduleRequest request) - throws CadenceError { + public CompletableFuture CreateSchedule(CreateScheduleRequest request) { return impl.CreateSchedule(request); } @Override - public DescribeScheduleResponse DescribeSchedule(DescribeScheduleRequest request) - throws CadenceError { + public CompletableFuture DescribeSchedule( + DescribeScheduleRequest request) { return impl.DescribeSchedule(request); } @Override - public UpdateScheduleResponse UpdateSchedule(UpdateScheduleRequest request) - throws CadenceError { + public CompletableFuture UpdateSchedule(UpdateScheduleRequest request) { return impl.UpdateSchedule(request); } @Override - public DeleteScheduleResponse DeleteSchedule(DeleteScheduleRequest request) - throws CadenceError { + public CompletableFuture DeleteSchedule(DeleteScheduleRequest request) { return impl.DeleteSchedule(request); } @Override - public PauseScheduleResponse PauseSchedule(PauseScheduleRequest request) throws CadenceError { + public CompletableFuture PauseSchedule(PauseScheduleRequest request) { return impl.PauseSchedule(request); } @Override - public UnpauseScheduleResponse UnpauseSchedule(UnpauseScheduleRequest request) - throws CadenceError { + public CompletableFuture UnpauseSchedule( + UnpauseScheduleRequest request) { return impl.UnpauseSchedule(request); } @Override - public BackfillScheduleResponse BackfillSchedule(BackfillScheduleRequest request) - throws CadenceError { + public CompletableFuture BackfillSchedule( + BackfillScheduleRequest request) { return impl.BackfillSchedule(request); } @Override - public ListSchedulesResponse ListSchedules(ListSchedulesRequest request) throws CadenceError { + public CompletableFuture ListSchedules(ListSchedulesRequest request) { return impl.ListSchedules(request); } diff --git a/src/main/java/com/uber/cadence/internal/sync/WorkflowClientInternal.java b/src/main/java/com/uber/cadence/internal/sync/WorkflowClientInternal.java index b19b87c4c..9ed10cb48 100644 --- a/src/main/java/com/uber/cadence/internal/sync/WorkflowClientInternal.java +++ b/src/main/java/com/uber/cadence/internal/sync/WorkflowClientInternal.java @@ -206,7 +206,7 @@ public ActivityCompletionClient newActivityCompletionClient() { @Override public ScheduleClient scheduleClient() { - throw new UnsupportedOperationException("not implemented"); + return new ScheduleClientImpl(workflowService, clientOptions.getDomain()); } @Override diff --git a/src/main/java/com/uber/cadence/internal/testservice/TestWorkflowService.java b/src/main/java/com/uber/cadence/internal/testservice/TestWorkflowService.java index 03360e23b..887bb3772 100644 --- a/src/main/java/com/uber/cadence/internal/testservice/TestWorkflowService.java +++ b/src/main/java/com/uber/cadence/internal/testservice/TestWorkflowService.java @@ -1269,45 +1269,45 @@ public void RefreshWorkflowTasks( } @Override - public CreateScheduleResponse CreateSchedule(CreateScheduleRequest request) throws CadenceError { + public CompletableFuture CreateSchedule(CreateScheduleRequest request) { throw new UnsupportedOperationException("not implemented"); } @Override - public DescribeScheduleResponse DescribeSchedule(DescribeScheduleRequest request) - throws CadenceError { + public CompletableFuture DescribeSchedule( + DescribeScheduleRequest request) { throw new UnsupportedOperationException("not implemented"); } @Override - public UpdateScheduleResponse UpdateSchedule(UpdateScheduleRequest request) throws CadenceError { + public CompletableFuture UpdateSchedule(UpdateScheduleRequest request) { throw new UnsupportedOperationException("not implemented"); } @Override - public DeleteScheduleResponse DeleteSchedule(DeleteScheduleRequest request) throws CadenceError { + public CompletableFuture DeleteSchedule(DeleteScheduleRequest request) { throw new UnsupportedOperationException("not implemented"); } @Override - public PauseScheduleResponse PauseSchedule(PauseScheduleRequest request) throws CadenceError { + public CompletableFuture PauseSchedule(PauseScheduleRequest request) { throw new UnsupportedOperationException("not implemented"); } @Override - public UnpauseScheduleResponse UnpauseSchedule(UnpauseScheduleRequest request) - throws CadenceError { + public CompletableFuture UnpauseSchedule( + UnpauseScheduleRequest request) { throw new UnsupportedOperationException("not implemented"); } @Override - public BackfillScheduleResponse BackfillSchedule(BackfillScheduleRequest request) - throws CadenceError { + public CompletableFuture BackfillSchedule( + BackfillScheduleRequest request) { throw new UnsupportedOperationException("not implemented"); } @Override - public ListSchedulesResponse ListSchedules(ListSchedulesRequest request) throws CadenceError { + public CompletableFuture ListSchedules(ListSchedulesRequest request) { throw new UnsupportedOperationException("not implemented"); } diff --git a/src/main/java/com/uber/cadence/serviceclient/IWorkflowService.java b/src/main/java/com/uber/cadence/serviceclient/IWorkflowService.java index 2f2e8f0be..f1cb223a0 100644 --- a/src/main/java/com/uber/cadence/serviceclient/IWorkflowService.java +++ b/src/main/java/com/uber/cadence/serviceclient/IWorkflowService.java @@ -817,88 +817,43 @@ void RefreshWorkflowTasks( RefreshWorkflowTasksRequest request, AsyncMethodCallback resultHandler) throws CadenceError; - /** - * Creates a new schedule in the domain. - * - * @throws BadRequestError if the schedule spec is invalid (e.g. bad cron expression) - * @throws WorkflowExecutionAlreadyStartedError if a schedule with the same ID already exists - * @throws EntityNotExistsError if the domain does not exist - * @throws CadenceError on other server errors - */ - CreateScheduleResponse CreateSchedule(CreateScheduleRequest request) - throws BadRequestError, WorkflowExecutionAlreadyStartedError, EntityNotExistsError, - ServiceBusyError, DomainNotActiveError, LimitExceededError, CadenceError; + /** Creates a new schedule in the domain. Failures are delivered via the returned future. */ + CompletableFuture CreateSchedule(CreateScheduleRequest request); /** - * Returns the full configuration and runtime state of a schedule. - * - * @throws EntityNotExistsError if the schedule does not exist - * @throws CadenceError on other server errors + * Returns the full configuration and runtime state of a schedule. Failures are delivered via the + * returned future. */ - DescribeScheduleResponse DescribeSchedule(DescribeScheduleRequest request) - throws BadRequestError, EntityNotExistsError, ServiceBusyError, CadenceError; + CompletableFuture DescribeSchedule(DescribeScheduleRequest request); /** * Replaces the spec, action, and policies of an existing schedule. * *

The server replaces the entire schedule configuration. Any field not included in the request * is zeroed out; always call describe first and resubmit the full configuration. - * - * @throws EntityNotExistsError if the schedule does not exist - * @throws CadenceError on other server errors */ - UpdateScheduleResponse UpdateSchedule(UpdateScheduleRequest request) - throws BadRequestError, EntityNotExistsError, ServiceBusyError, DomainNotActiveError, - LimitExceededError, CadenceError; + CompletableFuture UpdateSchedule(UpdateScheduleRequest request); /** * Permanently deletes a schedule. Running workflows triggered by this schedule are not affected. - * - * @throws EntityNotExistsError if the schedule does not exist - * @throws CadenceError on other server errors */ - DeleteScheduleResponse DeleteSchedule(DeleteScheduleRequest request) - throws BadRequestError, EntityNotExistsError, ServiceBusyError, DomainNotActiveError, - CadenceError; + CompletableFuture DeleteSchedule(DeleteScheduleRequest request); - /** - * Pauses a schedule. No new runs will be triggered while the schedule is paused. - * - * @throws EntityNotExistsError if the schedule does not exist - * @throws CadenceError on other server errors - */ - PauseScheduleResponse PauseSchedule(PauseScheduleRequest request) - throws BadRequestError, EntityNotExistsError, ServiceBusyError, DomainNotActiveError, - CadenceError; + /** Pauses a schedule. No new runs will be triggered while the schedule is paused. */ + CompletableFuture PauseSchedule(PauseScheduleRequest request); /** * Resumes a paused schedule. Optionally overrides the catch-up policy for the first pass on * resume. - * - * @throws EntityNotExistsError if the schedule does not exist - * @throws CadenceError on other server errors */ - UnpauseScheduleResponse UnpauseSchedule(UnpauseScheduleRequest request) - throws BadRequestError, EntityNotExistsError, ServiceBusyError, DomainNotActiveError, - CadenceError; + CompletableFuture UnpauseSchedule(UnpauseScheduleRequest request); /** * Triggers runs for all scheduled times within a historical time range. Backfill runs are subject * to the configured overlap policy. - * - * @throws EntityNotExistsError if the schedule does not exist - * @throws CadenceError on other server errors */ - BackfillScheduleResponse BackfillSchedule(BackfillScheduleRequest request) - throws BadRequestError, EntityNotExistsError, ServiceBusyError, DomainNotActiveError, - LimitExceededError, CadenceError; + CompletableFuture BackfillSchedule(BackfillScheduleRequest request); - /** - * Lists schedules in a domain, paginated. - * - * @throws EntityNotExistsError if the domain does not exist - * @throws CadenceError on other server errors - */ - ListSchedulesResponse ListSchedules(ListSchedulesRequest request) - throws BadRequestError, EntityNotExistsError, ServiceBusyError, CadenceError; + /** Lists schedules in a domain, paginated. */ + CompletableFuture ListSchedules(ListSchedulesRequest request); } diff --git a/src/main/java/com/uber/cadence/serviceclient/IWorkflowServiceBase.java b/src/main/java/com/uber/cadence/serviceclient/IWorkflowServiceBase.java index 4d08c67fd..ad2d3c232 100644 --- a/src/main/java/com/uber/cadence/serviceclient/IWorkflowServiceBase.java +++ b/src/main/java/com/uber/cadence/serviceclient/IWorkflowServiceBase.java @@ -724,45 +724,45 @@ public CompletableFuture isHealthy() { } @Override - public CreateScheduleResponse CreateSchedule(CreateScheduleRequest request) throws CadenceError { + public CompletableFuture CreateSchedule(CreateScheduleRequest request) { throw new UnsupportedOperationException("unimplemented"); } @Override - public DescribeScheduleResponse DescribeSchedule(DescribeScheduleRequest request) - throws CadenceError { + public CompletableFuture DescribeSchedule( + DescribeScheduleRequest request) { throw new UnsupportedOperationException("unimplemented"); } @Override - public UpdateScheduleResponse UpdateSchedule(UpdateScheduleRequest request) throws CadenceError { + public CompletableFuture UpdateSchedule(UpdateScheduleRequest request) { throw new UnsupportedOperationException("unimplemented"); } @Override - public DeleteScheduleResponse DeleteSchedule(DeleteScheduleRequest request) throws CadenceError { + public CompletableFuture DeleteSchedule(DeleteScheduleRequest request) { throw new UnsupportedOperationException("unimplemented"); } @Override - public PauseScheduleResponse PauseSchedule(PauseScheduleRequest request) throws CadenceError { + public CompletableFuture PauseSchedule(PauseScheduleRequest request) { throw new UnsupportedOperationException("unimplemented"); } @Override - public UnpauseScheduleResponse UnpauseSchedule(UnpauseScheduleRequest request) - throws CadenceError { + public CompletableFuture UnpauseSchedule( + UnpauseScheduleRequest request) { throw new UnsupportedOperationException("unimplemented"); } @Override - public BackfillScheduleResponse BackfillSchedule(BackfillScheduleRequest request) - throws CadenceError { + public CompletableFuture BackfillSchedule( + BackfillScheduleRequest request) { throw new UnsupportedOperationException("unimplemented"); } @Override - public ListSchedulesResponse ListSchedules(ListSchedulesRequest request) throws CadenceError { + public CompletableFuture ListSchedules(ListSchedulesRequest request) { throw new UnsupportedOperationException("unimplemented"); } } diff --git a/src/main/java/com/uber/cadence/serviceclient/WorkflowServiceGrpc.java b/src/main/java/com/uber/cadence/serviceclient/WorkflowServiceGrpc.java index c25a371e5..8294570ef 100644 --- a/src/main/java/com/uber/cadence/serviceclient/WorkflowServiceGrpc.java +++ b/src/main/java/com/uber/cadence/serviceclient/WorkflowServiceGrpc.java @@ -19,6 +19,7 @@ import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; import com.uber.cadence.*; import com.uber.cadence.internal.compatibility.proto.mappers.*; @@ -801,46 +802,110 @@ public void RefreshWorkflowTasks(RefreshWorkflowTasksRequest request) } @Override - public CreateScheduleResponse CreateSchedule(CreateScheduleRequest request) throws CadenceError { - throw new UnsupportedOperationException("not implemented"); + public CompletableFuture CreateSchedule(CreateScheduleRequest request) { + try { + return toCompletableFuture( + grpcServiceStubs + .scheduleFutureStub() + .createSchedule(RequestMapper.createScheduleRequest(request)), + ResponseMapper::createScheduleResponse); + } catch (Exception e) { + return failedFuture(e); + } } @Override - public DescribeScheduleResponse DescribeSchedule(DescribeScheduleRequest request) - throws CadenceError { - throw new UnsupportedOperationException("not implemented"); + public CompletableFuture DescribeSchedule( + DescribeScheduleRequest request) { + try { + return toCompletableFuture( + grpcServiceStubs + .scheduleFutureStub() + .describeSchedule(RequestMapper.describeScheduleRequest(request)), + ResponseMapper::describeScheduleResponse); + } catch (Exception e) { + return failedFuture(e); + } } @Override - public UpdateScheduleResponse UpdateSchedule(UpdateScheduleRequest request) throws CadenceError { - throw new UnsupportedOperationException("not implemented"); + public CompletableFuture UpdateSchedule(UpdateScheduleRequest request) { + try { + return toCompletableFuture( + grpcServiceStubs + .scheduleFutureStub() + .updateSchedule(RequestMapper.updateScheduleRequest(request)), + ResponseMapper::updateScheduleResponse); + } catch (Exception e) { + return failedFuture(e); + } } @Override - public DeleteScheduleResponse DeleteSchedule(DeleteScheduleRequest request) throws CadenceError { - throw new UnsupportedOperationException("not implemented"); + public CompletableFuture DeleteSchedule(DeleteScheduleRequest request) { + try { + return toCompletableFuture( + grpcServiceStubs + .scheduleFutureStub() + .deleteSchedule(RequestMapper.deleteScheduleRequest(request)), + ResponseMapper::deleteScheduleResponse); + } catch (Exception e) { + return failedFuture(e); + } } @Override - public PauseScheduleResponse PauseSchedule(PauseScheduleRequest request) throws CadenceError { - throw new UnsupportedOperationException("not implemented"); + public CompletableFuture PauseSchedule(PauseScheduleRequest request) { + try { + return toCompletableFuture( + grpcServiceStubs + .scheduleFutureStub() + .pauseSchedule(RequestMapper.pauseScheduleRequest(request)), + ResponseMapper::pauseScheduleResponse); + } catch (Exception e) { + return failedFuture(e); + } } @Override - public UnpauseScheduleResponse UnpauseSchedule(UnpauseScheduleRequest request) - throws CadenceError { - throw new UnsupportedOperationException("not implemented"); + public CompletableFuture UnpauseSchedule( + UnpauseScheduleRequest request) { + try { + return toCompletableFuture( + grpcServiceStubs + .scheduleFutureStub() + .unpauseSchedule(RequestMapper.unpauseScheduleRequest(request)), + ResponseMapper::unpauseScheduleResponse); + } catch (Exception e) { + return failedFuture(e); + } } @Override - public BackfillScheduleResponse BackfillSchedule(BackfillScheduleRequest request) - throws CadenceError { - throw new UnsupportedOperationException("not implemented"); + public CompletableFuture BackfillSchedule( + BackfillScheduleRequest request) { + try { + return toCompletableFuture( + grpcServiceStubs + .scheduleFutureStub() + .backfillSchedule(RequestMapper.backfillScheduleRequest(request)), + ResponseMapper::backfillScheduleResponse); + } catch (Exception e) { + return failedFuture(e); + } } @Override - public ListSchedulesResponse ListSchedules(ListSchedulesRequest request) throws CadenceError { - throw new UnsupportedOperationException("not implemented"); + public CompletableFuture ListSchedules(ListSchedulesRequest request) { + try { + return toCompletableFuture( + grpcServiceStubs + .scheduleFutureStub() + .listSchedules(RequestMapper.listSchedulesRequest(request)), + ResponseMapper::listSchedulesResponse); + } catch (Exception e) { + return failedFuture(e); + } } @Override @@ -1414,6 +1479,32 @@ private CadenceError toServiceClientException(Throwable t) { } } + private CompletableFuture toCompletableFuture( + ListenableFuture

listenableFuture, Function mapper) { + CompletableFuture future = new CompletableFuture<>(); + Futures.addCallback( + listenableFuture, + new FutureCallback

() { + @Override + public void onSuccess(P result) { + future.complete(mapper.apply(result)); + } + + @Override + public void onFailure(Throwable t) { + future.completeExceptionally(toServiceClientException(t)); + } + }, + executor); + return future; + } + + private CompletableFuture failedFuture(Exception e) { + CompletableFuture future = new CompletableFuture<>(); + future.completeExceptionally(toServiceClientException(e)); + return future; + } + private FutureCallback toFutureCallback( AsyncMethodCallback resultHandler, Function mapper) { return new FutureCallback() {