Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
43 commits
Select commit Hold shift + click to select a range
9400ace
annotation support for system task workers
v1r3n Jan 25, 2026
0a7580b
support TaskContext
v1r3n Jan 28, 2026
9fcf5eb
Fix: Restore missing postpone logic from old sweeper (PR #703 regress…
manan164 Jan 23, 2026
df2ba8a
Refactor: Move sweeper to common package for Orkes compatibility
manan164 Jan 23, 2026
0caef27
Add tests for ExecutorUtils.computePostpone
manan164 Jan 28, 2026
8699b91
revert
manan164 Jan 28, 2026
a03f4ff
Add tests for ExecutorUtils.computePostpone
manan164 Jan 28, 2026
ee3596d
Fix TestWorkflowSweeper for new WorkflowSweeper API
manan164 Jan 28, 2026
d1aee10
revert
manan164 Jan 28, 2026
9e8e7b3
spotless.
manan164 Jan 28, 2026
5d5ade9
Move Orkes WorkflowSweeper from org.conductoross to com.netflix.condu…
manan164 Jan 28, 2026
e174e2c
spotless
manan164 Jan 29, 2026
85783d9
changes
manan164 Jan 29, 2026
58628ee
Fix: Update WorkflowSweeper import in AbstractSpecification
manan164 Jan 29, 2026
bf0f1f7
revert
manan164 Jan 29, 2026
537c4dd
Add Orkes improvements to WorkflowSweeper (without Orkes-specific code)
manan164 Jan 29, 2026
efd3768
Revert "Move Orkes WorkflowSweeper from org.conductoross to com.netfl…
manan164 Jan 29, 2026
17c92cf
Restore old WorkflowSweeper in org.conductoross package
manan164 Jan 29, 2026
a29a407
Restore ExecutorUtils to org.conductoross package
manan164 Jan 29, 2026
3f1d02a
simpler sweeper.
manan164 Jan 29, 2026
6eba38a
changes
manan164 Jan 29, 2026
6445e3b
Revert changes to com.netflix WorkflowSweeper - keep it disabled as-is
manan164 Jan 29, 2026
e558781
Remove ExecutorUtils from com.netflix package - only needed in org.co…
manan164 Jan 29, 2026
5a4de9a
Fix TestExecutorUtils to import ExecutorUtils from org.conductoross p…
manan164 Jan 29, 2026
3b4c88c
Remove ExecutorUtils - not needed for OrkesWorkflowSweeper parity
manan164 Jan 29, 2026
7f20847
Remove unused sweeper-related code: SweeperProperties in com.netflix,…
manan164 Jan 29, 2026
d18d454
Remove all unrelated changes - keep only sweeper changes
manan164 Jan 29, 2026
b93eab4
Restore ExecutorUtils - needed by WorkflowExecutorOps
manan164 Jan 29, 2026
454fb8f
Add explicit lock management to WorkflowSweeper
manan164 Jan 30, 2026
2b46e47
Fix OpenSearchProperties Environment injection for integration tests
manan164 Jan 30, 2026
d015be4
Fix WorkflowSweeper imports in test base classes
manan164 Jan 30, 2026
81bd987
spotless
manan164 Jan 30, 2026
d67357e
remove async annotation
manan164 Jan 30, 2026
fc97206
spotless
manan164 Jan 30, 2026
dd8ccaf
Merge branch 'main' into fix/restore-postpone-logic-from-old-sweeper
manan164 Jan 30, 2026
6295cb6
Merge pull request #728 from conductor-oss/fix/restore-postpone-logic…
v1r3n Jan 30, 2026
fc888f0
Keep legacy sweeper
manan164 Jan 30, 2026
e789f81
spotless
manan164 Jan 30, 2026
36b52d8
fix
manan164 Jan 30, 2026
1ae620a
Merge pull request #744 from conductor-oss/fix/keep-legacy-sweeper
manan164 Jan 30, 2026
befa983
removed duplicated classes
v1r3n Jan 31, 2026
68794a0
fixes and clean up
v1r3n Jan 31, 2026
0f02c10
Merge pull request #737 from conductor-oss/ai_orchestration
v1r3n Jan 31, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -105,15 +105,23 @@ public ExecutorService executorService(ConductorProperties conductorProperties)
@Bean
@Qualifier("taskMappersByTaskType")
public Map<String, TaskMapper> getTaskMappers(List<TaskMapper> taskMappers) {
return taskMappers.stream().collect(Collectors.toMap(TaskMapper::getTaskType, identity()));
// Return mutable map so annotated task mappers can be added
return taskMappers.stream()
.collect(
Collectors.toMap(
TaskMapper::getTaskType,
identity(),
(a, b) -> a,
java.util.HashMap::new));
}

@Bean
@Qualifier(ASYNC_SYSTEM_TASKS_QUALIFIER)
public Set<WorkflowSystemTask> asyncSystemTasks(Set<WorkflowSystemTask> allSystemTasks) {
// Return mutable set so annotated tasks can be added
return allSystemTasks.stream()
.filter(WorkflowSystemTask::isAsync)
.collect(Collectors.toUnmodifiableSet());
.collect(Collectors.toCollection(java.util.HashSet::new));
}

@Bean
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,7 @@ public DeciderService(
ExternalPayloadStorageUtils externalPayloadStorageUtils,
SystemTaskRegistry systemTaskRegistry,
@Qualifier("taskMappersByTaskType") Map<String, TaskMapper> taskMappers,
@Qualifier("annotatedTaskSystems") Map<String, TaskMapper> annotatedTaskSystems,
@Value("${conductor.app.taskPendingTimeThreshold:60m}")
Duration taskPendingTimeThreshold) {
this.idGenerator = idGenerator;
Expand All @@ -84,6 +85,19 @@ public DeciderService(
this.externalPayloadStorageUtils = externalPayloadStorageUtils;
this.taskPendingTimeThresholdMins = taskPendingTimeThreshold.toMinutes();
this.systemTaskRegistry = systemTaskRegistry;
LOGGER.info("taskMappers: {}", taskMappers.keySet());
LOGGER.info("annotatedTaskMappers: {}", annotatedTaskSystems.keySet());
// Add annotated mappers only if no existing mapper is registered for that task
// type
// Existing mappers take precedence over annotated ones
annotatedTaskSystems.forEach(
(taskType, mapper) -> {
if (this.taskMappers.putIfAbsent(taskType, mapper) != null) {
LOGGER.info(
"Skipping annotated mapper for '{}' - existing mapper already registered",
taskType);
}
});
}

public DeciderOutcome decide(WorkflowModel workflow) throws TerminateWorkflowException {
Expand Down Expand Up @@ -135,8 +149,10 @@ private DeciderOutcome decide(final WorkflowModel workflow, List<TaskModel> preS
boolean hasSuccessfulTerminateTask = false;
for (TaskModel task : workflow.getTasks()) {

// Filter the list of tasks and include only tasks that are not retried, not executed
// marked to be skipped and not part of System tasks that is DECISION, FORK, JOIN
// Filter the list of tasks and include only tasks that are not retried, not
// executed
// marked to be skipped and not part of System tasks that is DECISION, FORK,
// JOIN
// This list will be empty for a new workflow being started
if (!task.isRetried() && !task.getStatus().equals(SKIPPED) && !task.isExecuted()) {
pendingTasks.add(task);
Expand Down Expand Up @@ -186,7 +202,8 @@ private DeciderOutcome decide(final WorkflowModel workflow, List<TaskModel> preS
if (taskDefinition.isPresent()) {
checkTaskTimeout(taskDefinition.get(), pendingTask);
checkTaskPollTimeout(taskDefinition.get(), pendingTask);
// If the task has not been updated for "responseTimeoutSeconds" then mark task as
// If the task has not been updated for "responseTimeoutSeconds" then mark task
// as
// TIMED_OUT
if (isResponseTimedOut(taskDefinition.get(), pendingTask)) {
timeoutTask(taskDefinition.get(), pendingTask);
Expand Down Expand Up @@ -448,7 +465,8 @@ public boolean checkForWorkflowCompletion(final WorkflowModel workflow)
return false;
}

// If there is a TERMINATE task that has been executed successfuly then the workflow
// If there is a TERMINATE task that has been executed successfuly then the
// workflow
// should be marked as completed.
if (TERMINATE.name().equals(task.getTaskType())
&& task.getStatus().isTerminal()
Expand Down Expand Up @@ -869,7 +887,7 @@ public List<TaskModel> getTasksToBeScheduled(

String type = taskToSchedule.getType();

// get tasks already scheduled (in progress/terminal) for this workflow instance
// get tasks already scheduled (in progress/terminal) for this workflow instance
List<String> tasksInWorkflow =
workflow.getTasks().stream()
.filter(
Expand All @@ -892,10 +910,13 @@ public List<TaskModel> getTasksToBeScheduled(
.withDeciderService(this)
.build();

// For static forks, each branch of the fork creates a join task upon completion for
// dynamic forks, a join task is created with the fork and also with each branch of the
// For static forks, each branch of the fork creates a join task upon completion
// for
// dynamic forks, a join task is created with the fork and also with each branch
// of the
// fork.
// A new task must only be scheduled if a task, with the same reference name is not already
// A new task must only be scheduled if a task, with the same reference name is
// not already
// in this workflow instance
return taskMappers
.getOrDefault(type, taskMappers.get(USER_DEFINED.name()))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,13 +18,15 @@
import java.util.function.Function;
import java.util.stream.Collectors;

import org.springframework.context.annotation.DependsOn;
import org.springframework.stereotype.Component;

/**
* A container class that holds a mapping of system task types {@link
* com.netflix.conductor.common.metadata.tasks.TaskType} to {@link WorkflowSystemTask} instances.
*/
@Component
@DependsOn("workerTaskAnnotationScanner")
public class SystemTaskRegistry {

public static final String ASYNC_SYSTEM_TASKS_QUALIFIER = "asyncSystemTasks";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,9 @@

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.scheduling.annotation.Async;
import org.springframework.stereotype.Component;

import com.netflix.conductor.annotations.VisibleForTesting;
import com.netflix.conductor.common.metadata.tasks.TaskDef;
Expand All @@ -39,7 +41,13 @@
import static com.netflix.conductor.core.config.SchedulerConfiguration.SWEEPER_EXECUTOR_NAME;
import static com.netflix.conductor.core.utils.Utils.DECIDER_QUEUE;

// @Component
// Deprecated in favor of org.conductoross.conductor.core.execution.WorkflowSweeper
@Deprecated(forRemoval = true)
@Component
@ConditionalOnProperty(
name = "conductor.app.legacy.sweeper.enabled",
havingValue = "true",
matchIfMissing = false)
public class WorkflowSweeper {

private static final Logger LOGGER = LoggerFactory.getLogger(WorkflowSweeper.class);
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
/*
* Copyright 2021 Conductor Authors.
* <p>
* Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
* <p>
* http://www.apache.org/licenses/LICENSE-2.0
* <p>
* Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on
* an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the
* specific language governing permissions and limitations under the License.
*/
package com.netflix.conductor.sdk.workflow.task;

import java.lang.annotation.ElementType;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;

@Retention(RetentionPolicy.RUNTIME)
@Target(ElementType.PARAMETER)
public @interface InputParam {
String value();

boolean required() default false;
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
/*
* Copyright 2021 Conductor Authors.
* <p>
* Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
* <p>
* http://www.apache.org/licenses/LICENSE-2.0
* <p>
* Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on
* an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the
* specific language governing permissions and limitations under the License.
*/
package com.netflix.conductor.sdk.workflow.task;

import java.lang.annotation.ElementType;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;

@Retention(RetentionPolicy.RUNTIME)
@Target(ElementType.TYPE_USE)
public @interface OutputParam {
String value();
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
/*
* Copyright 2021 Conductor Authors.
* <p>
* Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
* <p>
* http://www.apache.org/licenses/LICENSE-2.0
* <p>
* Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on
* an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the
* specific language governing permissions and limitations under the License.
*/
package com.netflix.conductor.sdk.workflow.task;

import java.lang.annotation.ElementType;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;

/** Identifies a simple worker task. */
@Retention(RetentionPolicy.RUNTIME)
@Target({ElementType.METHOD})
public @interface WorkerTask {
String value();

// No. of threads to use for executing the task
int threadCount() default 1;

int pollingInterval() default 100;

String domain() default "";

// In millis
int pollTimeout() default 100;

// number of task pollers
// default is 1 which is good enough for most use cases
// a number higher than 1 will have concurrent pollers doing poll and execute
int pollerCount() default 1;
}
Loading
Loading