Skip to content
Open
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
29 changes: 29 additions & 0 deletions Conductor/Client/Extensions/DependencyInjectionExtensions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,9 @@
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using System;
using System.Linq;
using System.Net.Http;
using System.Reflection;

namespace Conductor.Client.Extensions
{
Expand Down Expand Up @@ -63,5 +65,32 @@ public static IServiceCollection WithHostedService(this IServiceCollection servi
services.AddHostedService<WorkflowTaskService>();
return services;
}

public static IServiceCollection ConfigureConductorWorkerDiscovery(this IServiceCollection services, Action<WorkerDiscoveryOptions> configure)
{
if (configure == null)
{
throw new ArgumentNullException(nameof(configure));
}

services.AddOptions<WorkerDiscoveryOptions>().Configure(configure);

return services;
}

public static IServiceCollection ConfigureConductorWorkerDiscovery(this IServiceCollection services, WorkerDiscoveryOptions discoveryOptions)
{
if (discoveryOptions == null)
{
throw new ArgumentNullException(nameof(discoveryOptions));
}

var assemblies = discoveryOptions.Assemblies?.ToArray() ?? Array.Empty<Assembly>();

return services.ConfigureConductorWorkerDiscovery(options =>
{
options.Assemblies = assemblies;
});
}
}
}
17 changes: 17 additions & 0 deletions Conductor/Client/Worker/WorkerDiscoveryOptions.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
using System;
using System.Collections.Generic;
using System.Reflection;

namespace Conductor.Client.Worker
{

public sealed class WorkerDiscoveryOptions
{
public bool EnableAttributeDiscovery { get; set; } = true;

/// <summary>
/// Empty means: preserve current behaviour and scan all loaded assemblies.
/// </summary>
public IReadOnlyCollection<Assembly> Assemblies { get; set; } = Array.Empty<Assembly>();
}
}
24 changes: 21 additions & 3 deletions Conductor/Client/Worker/WorkflowTaskCoordinator.cs
Original file line number Diff line number Diff line change
Expand Up @@ -13,8 +13,10 @@
using Conductor.Client.Interfaces;
using Conductor.Client.Telemetry;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
using System;
using System.Collections.Generic;
using System.Linq;
using System.Reflection;
using System.Threading;
using System.Threading.Tasks;
Expand All @@ -29,15 +31,23 @@ internal class WorkflowTaskCoordinator : IWorkflowTaskCoordinator
private readonly HashSet<IWorkflowTaskExecutor> _workers;
private readonly IWorkflowTaskClient _client;
private readonly MetricsCollector _metrics;
private readonly WorkerDiscoveryOptions _workerDiscoveryOptions;

public WorkflowTaskCoordinator(IWorkflowTaskClient client, ILogger<WorkflowTaskCoordinator> logger, ILogger<WorkflowTaskExecutor> loggerWorkflowTaskExecutor, ILogger<WorkflowTaskMonitor> loggerWorkflowTaskMonitor, MetricsCollector metrics = null)
public WorkflowTaskCoordinator(IWorkflowTaskClient client,
ILogger<WorkflowTaskCoordinator> logger,
ILogger<WorkflowTaskExecutor> loggerWorkflowTaskExecutor,
ILogger<WorkflowTaskMonitor> loggerWorkflowTaskMonitor,
MetricsCollector metrics = null,
IOptions<WorkerDiscoveryOptions> workerDiscoveryOptions = null)
{
_logger = logger;
_client = client;
_workers = new HashSet<IWorkflowTaskExecutor>();
_loggerWorkflowTaskExecutor = loggerWorkflowTaskExecutor;
_loggerWorkflowTaskMonitor = loggerWorkflowTaskMonitor;
_metrics = metrics;

_workerDiscoveryOptions = workerDiscoveryOptions?.Value ?? new WorkerDiscoveryOptions();
}

public async Task Start(CancellationToken token)
Expand All @@ -46,7 +56,12 @@ public async Task Start(CancellationToken token)
token.ThrowIfCancellationRequested();

_logger.LogDebug("Starting workers...");
DiscoverWorkers();

if (_workerDiscoveryOptions.EnableAttributeDiscovery)
{
DiscoverWorkers();
}

var runningWorkers = new List<Task>();
foreach (var worker in _workers)
{
Expand All @@ -72,14 +87,17 @@ public void RegisterWorker(IWorkflowTask worker)

private void DiscoverWorkers()
{
foreach (var assembly in AppDomain.CurrentDomain.GetAssemblies())
var assemblies = _workerDiscoveryOptions.Assemblies?.Any() == true ? _workerDiscoveryOptions.Assemblies : AppDomain.CurrentDomain.GetAssemblies();

foreach (var assembly in assemblies)
{
foreach (var type in assembly.GetTypes())
{
if (type.GetCustomAttribute<WorkerTask>() == null)
{
continue;
}

foreach (var method in type.GetMethods())
{
var workerTask = method.GetCustomAttribute<WorkerTask>();
Expand Down
67 changes: 66 additions & 1 deletion Tests/Extensions/DependencyInjectionExtensionsTests.cs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
/*
/*
* Copyright 2024 Conductor Authors.
* <p>
* Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with
Expand All @@ -12,7 +12,11 @@
*/
using Conductor.Client.Extensions;
using Conductor.Client.Telemetry;
using Conductor.Client.Worker;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Options;
using System.Collections.Generic;
using System.Reflection;
using Xunit;

namespace Tests.Extensions
Expand Down Expand Up @@ -88,5 +92,66 @@ public void AddConductorWorker_WithNullConfiguration_CreatesDefault()
var provider = services.BuildServiceProvider();
Assert.NotNull(provider.GetService<Conductor.Client.Configuration>());
}

[Fact]
public void ConfigureConductorWorkerDiscovery_WithoutConfiguration_UsesEmptyAssemblyCollection()
{
var services = new ServiceCollection();

services.AddConductorWorker();

var provider = services.BuildServiceProvider();
var options = provider.GetRequiredService<IOptions<WorkerDiscoveryOptions>>().Value;

Assert.Empty(options.Assemblies);
}

[Fact]
public void ConfigureConductorWorkerDiscovery_WithAction_ConfiguresAssemblyCollection()
{
var services = new ServiceCollection();
var assembly = typeof(DependencyInjectionExtensionsTests).Assembly;

services.ConfigureConductorWorkerDiscovery(options =>
options.Assemblies = new[] { assembly });

var provider = services.BuildServiceProvider();
var options = provider.GetRequiredService<IOptions<WorkerDiscoveryOptions>>().Value;

Assert.Contains(assembly, options.Assemblies);
}

[Fact]
public void ConfigureConductorWorkerDiscovery_WithAttributeDiscoveryDisabled_ConfiguresOptions()
{
var services = new ServiceCollection();

services.ConfigureConductorWorkerDiscovery(options =>
options.EnableAttributeDiscovery = false);

var provider = services.BuildServiceProvider();
var options = provider.GetRequiredService<IOptions<WorkerDiscoveryOptions>>().Value;

Assert.False(options.EnableAttributeDiscovery);
}

[Fact]
public void ConfigureConductorWorkerDiscovery_WithOptions_CopiesAssemblyCollection()
{
var services = new ServiceCollection();
var assembly = typeof(DependencyInjectionExtensionsTests).Assembly;
var assemblies = new List<Assembly> { assembly };

services.ConfigureConductorWorkerDiscovery(new WorkerDiscoveryOptions
{
Assemblies = assemblies
});
assemblies.Clear();

var provider = services.BuildServiceProvider();
var options = provider.GetRequiredService<IOptions<WorkerDiscoveryOptions>>().Value;

Assert.Contains(assembly, options.Assemblies);
}
}
}
27 changes: 27 additions & 0 deletions docs/workers.md
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,33 @@ settings. See [deployment-scaling.md](deployment-scaling.md) for sizing guidance
`IServiceCollection`, so workers can take constructor dependencies and participate in
the host's lifetime. This is the preferred shape for anything beyond a sample.

### Choose one registration path per worker

The SDK supports two worker registration paths:

- Register an `IWorkflowTask` explicitly in the service collection, for example with
`AddConductorWorkflowTask` or `ServiceDescriptor.Singleton<IWorkflowTask, TWorker>()`.
- Discover methods annotated with `[WorkerTask]` when the worker host starts.

Choose one path for each worker. Do not explicitly register a worker that is also
discovered through `[WorkerTask]`, or the host creates two polling workers for it.

By default, attribute discovery scans all assemblies loaded in the process. To limit
the scan to the assembly containing your annotated workers, configure discovery while
building the service collection:

```csharp
services.AddConductorWorker(configuration);
services.ConfigureConductorWorkerDiscovery(options =>
{
options.Assemblies = new[] { typeof(MyAnnotatedWorker).Assembly };
});
services.WithHostedService();
```

This setting affects only `[WorkerTask]` discovery. It does not affect any
`IWorkflowTask` registered in the service collection.

## Metrics

The worker framework records polling, execution, update, and error metrics via
Expand Down
Loading