diff --git a/Conductor/Client/Worker/WorkflowTaskExecutor.cs b/Conductor/Client/Worker/WorkflowTaskExecutor.cs index c081f775..2fa1d08a 100644 --- a/Conductor/Client/Worker/WorkflowTaskExecutor.cs +++ b/Conductor/Client/Worker/WorkflowTaskExecutor.cs @@ -60,6 +60,7 @@ public WorkflowTaskExecutor( _logger = logger; _taskClient = client; _worker = worker; + _workerSettings = workflowTaskConfiguration ?? worker.WorkerSettings; _workflowTaskMonitor = workflowTaskMonitor; _metrics = metrics; } diff --git a/Tests/Worker/WorkflowTaskExecutorTests.cs b/Tests/Worker/WorkflowTaskExecutorTests.cs index 5c5ebd2f..fd4a4c1d 100644 --- a/Tests/Worker/WorkflowTaskExecutorTests.cs +++ b/Tests/Worker/WorkflowTaskExecutorTests.cs @@ -126,6 +126,72 @@ 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()); + + 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.Instance); + + var executor = new WorkflowTaskExecutor(NullLogger.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()); + + var worker = new FakeWorker(taskType, batchSize); + WorkflowTaskExecutorConfiguration config = null; + var monitor = new WorkflowTaskMonitor(NullLogger.Instance); + + var executor = new WorkflowTaskExecutor(NullLogger.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", @@ -133,9 +199,11 @@ private WorkflowTaskExecutor CreateExecutor( { var worker = new FakeWorker(taskType, batchSize); var monitor = new WorkflowTaskMonitor(NullLogger.Instance); - return new WorkflowTaskExecutor( - NullLogger.Instance, - taskClient, worker, monitor, _metrics); + return new WorkflowTaskExecutor(NullLogger.Instance, + taskClient, + worker, + monitor, + _metrics); } private static async Task RunOnceAndWait(WorkflowTaskExecutor executor) @@ -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 returnTasks = null, Exception pollException = null, @@ -194,6 +266,10 @@ public string UpdateTask(TaskResult result) public Task> PollTaskAsync(string taskType, string workerId, string domain, int count) { + LastWorkerId = workerId; + LastDomain = domain; + LastRequestedTaskCount = count; + return Task.FromResult(PollTask(taskType, workerId, domain, count)); }