Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
1 change: 1 addition & 0 deletions Conductor/Client/Worker/WorkflowTaskExecutor.cs
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,7 @@ public WorkflowTaskExecutor(
_logger = logger;
_taskClient = client;
_worker = worker;
_workerSettings = workflowTaskConfiguration ?? worker.WorkerSettings;
_workflowTaskMonitor = workflowTaskMonitor;
_metrics = metrics;
}
Expand Down
82 changes: 79 additions & 3 deletions Tests/Worker/WorkflowTaskExecutorTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -126,16 +126,84 @@ public async Task CancellationToken_BreaksWorkLoop()
Assert.Equal(runTask, completed);
}

[Fact]
public async Task CreateExecutor_WithExecutorConfig()
{
// Arrange
string taskType = "test_task";
int batchSize = 10;

var taskClient = new FakeTaskClient(returnTasks: new List<Conductor.Client.Models.Task>());

var worker = new FakeWorker(taskType, batchSize);
var config = new WorkflowTaskExecutorConfiguration
{
WorkerId = "configured-worker",
Domain = "configured-domain",
BatchSize = 3,
PollInterval = TimeSpan.FromMilliseconds(10)
};
var monitor = new WorkflowTaskMonitor(NullLogger<WorkflowTaskMonitor>.Instance);

var executor = new WorkflowTaskExecutor(NullLogger<WorkflowTaskExecutor>.Instance,
taskClient,
worker,
config,
monitor,
_metrics);

// Act
await RunOnceAndWait(executor);

// Assert

Assert.Equal("configured-worker", taskClient.LastWorkerId);
Assert.Equal("configured-domain", taskClient.LastDomain);
Assert.Equal(3, taskClient.LastRequestedTaskCount);
}

[Fact]
public async Task CreateExecutor_WithWorkerConfig()
{
// Arrange
string taskType = "test_task";
int batchSize = 10;

var taskClient = new FakeTaskClient(returnTasks: new List<Conductor.Client.Models.Task>());

var worker = new FakeWorker(taskType, batchSize);
WorkflowTaskExecutorConfiguration config = null;
var monitor = new WorkflowTaskMonitor(NullLogger<WorkflowTaskMonitor>.Instance);

var executor = new WorkflowTaskExecutor(NullLogger<WorkflowTaskExecutor>.Instance,
taskClient,
worker,
config,
monitor,
_metrics);

// Act
await RunOnceAndWait(executor);

// Assert

Assert.Equal("test-worker-1", taskClient.LastWorkerId);
Assert.Equal("test", taskClient.LastDomain);
Assert.Equal(10, taskClient.LastRequestedTaskCount);
}

private WorkflowTaskExecutor CreateExecutor(
IWorkflowTaskClient taskClient,
string taskType = "test_task",
int batchSize = 10)
{
var worker = new FakeWorker(taskType, batchSize);
var monitor = new WorkflowTaskMonitor(NullLogger<WorkflowTaskMonitor>.Instance);
return new WorkflowTaskExecutor(
NullLogger<WorkflowTaskExecutor>.Instance,
taskClient, worker, monitor, _metrics);
return new WorkflowTaskExecutor(NullLogger<WorkflowTaskExecutor>.Instance,
taskClient,
worker,
monitor,
_metrics);
}

private static async Task RunOnceAndWait(WorkflowTaskExecutor executor)
Expand Down Expand Up @@ -168,6 +236,10 @@ private class FakeTaskClient : IWorkflowTaskClient
public int PollCount { get; private set; }
public int UpdateCount { get; private set; }

public string LastWorkerId { get; private set; }
public string LastDomain { get; private set; }
public int LastRequestedTaskCount { get; private set; }

public FakeTaskClient(
List<Conductor.Client.Models.Task> returnTasks = null,
Exception pollException = null,
Expand All @@ -194,6 +266,10 @@ public string UpdateTask(TaskResult result)

public Task<List<Conductor.Client.Models.Task>> PollTaskAsync(string taskType, string workerId, string domain, int count)
{
LastWorkerId = workerId;
LastDomain = domain;
LastRequestedTaskCount = count;

return Task.FromResult(PollTask(taskType, workerId, domain, count));
}

Expand Down
Loading