diff --git a/docs/postgresql/ado-net.md b/docs/postgresql/ado-net.md index b644484..6653b16 100644 --- a/docs/postgresql/ado-net.md +++ b/docs/postgresql/ado-net.md @@ -35,6 +35,18 @@ services.UsePostgreSqlAdoNetOutbox(options => }); ``` +## Worker Connections + +`UseWorkerConnectionFactory(...)` is required when the application processes +the outbox, either through `ITinyOutboxProcessor` directly or through the +hosted worker. + +The hosted worker validates that the factory is configured before polling. It +does not invoke the delegate or open a database connection during startup +validation. Missing configuration stops startup; failures creating or opening +a configured connection occur during processing and follow the worker's +operational retry behavior. + ## Publishing Transaction Ownership TinyEvents: diff --git a/docs/postgresql/ef-core.md b/docs/postgresql/ef-core.md index eb02b16..e908701 100644 --- a/docs/postgresql/ef-core.md +++ b/docs/postgresql/ef-core.md @@ -57,6 +57,17 @@ The provider option controls SQL claiming and marking. The model builder extensi The PostgreSQL mapping uses `text` for `EventType`, `Payload`, `ClaimedBy`, and `LastError`. +## Worker Startup Validation + +The hosted worker validates that `TDbContext` uses the Npgsql EF Core provider +before polling. A missing or different EF Core provider stops startup because +the worker store executes PostgreSQL-specific claim and mark commands. + +Validation reads EF Core provider metadata only. It does not open a database +connection, check credentials, inspect the schema, or run migrations. Connection +failures after startup remain operational iteration failures and follow the +worker's retry behavior. + ## Worker Claiming The PostgreSQL EF Core store opens the underlying relational connection when needed and executes PostgreSQL claim/mark statements. diff --git a/docs/sql-server/ado-net.md b/docs/sql-server/ado-net.md index de27ac3..708bb86 100644 --- a/docs/sql-server/ado-net.md +++ b/docs/sql-server/ado-net.md @@ -35,6 +35,18 @@ services.UseSqlServerAdoNetOutbox(options => }); ``` +## Worker Connections + +`UseWorkerConnectionFactory(...)` is required when the application processes +the outbox, either through `ITinyOutboxProcessor` directly or through the +hosted worker. + +The hosted worker validates that the factory is configured before polling. It +does not invoke the delegate or open a database connection during startup +validation. Missing configuration stops startup; failures creating or opening +a configured connection occur during processing and follow the worker's +operational retry behavior. + ## Publishing Transaction Ownership TinyEvents: diff --git a/docs/sql-server/ef-core.md b/docs/sql-server/ef-core.md index 044879e..276ec4a 100644 --- a/docs/sql-server/ef-core.md +++ b/docs/sql-server/ef-core.md @@ -57,6 +57,17 @@ The provider option controls SQL claiming and marking. The model builder extensi The SQL Server mapping uses `NVARCHAR(512)` for `EventType`, `NVARCHAR(MAX)` for `Payload`, `NVARCHAR(256)` for `ClaimedBy`, and `NVARCHAR(MAX)` for `LastError`. +## Worker Startup Validation + +The hosted worker validates that `TDbContext` uses the SQL Server EF Core +provider before polling. A missing or different EF Core provider stops startup +because the worker store executes SQL Server-specific claim and mark commands. + +Validation reads EF Core provider metadata only. It does not open a database +connection, check credentials, inspect the schema, or run migrations. Connection +failures after startup remain operational iteration failures and follow the +worker's retry behavior. + ## Worker Claiming The SQL Server EF Core store opens the underlying relational connection when needed and executes SQL Server claim/mark statements. diff --git a/docs/workers.md b/docs/workers.md index 612d443..472bf71 100644 --- a/docs/workers.md +++ b/docs/workers.md @@ -88,6 +88,24 @@ On shutdown, TinyEvents does not scan and release claims. If processing does not A hosted worker can remain running while processing iterations repeatedly fail, for example during a database outage. Treat worker logs and host-level health checks as part of production operations. +## Startup Validation + +Before polling, the hosted worker resolves the processing graph once in a +temporary scope. This validates the processor dependencies and generated +dispatcher registrations without claiming or processing messages. + +A startup validation failure escapes to the host. It is not logged, counted, +or retried as a processing-iteration failure because changing the application +configuration or registrations is required to recover. + +After validation succeeds, operational iteration failures are logged and the +worker continues polling. Consumer failures, deserialization failures, unknown +event types stored in messages, and lease loss remain message-level outcomes +handled by the processor. + +Database schema initialization is separate from runtime graph validation. +TinyEvents does not currently run schema migrations as part of worker startup. + ## Runtime Logging TinyEvents uses `Microsoft.Extensions.Logging`. The application owns log diff --git a/src/TinyEvents.PostgreSql.AdoNet/Connections/TinyPostgreSqlAdoNetWorkerConnectionFactory.cs b/src/TinyEvents.PostgreSql.AdoNet/Connections/TinyPostgreSqlAdoNetWorkerConnectionFactory.cs index ce4d73e..0b615fa 100644 --- a/src/TinyEvents.PostgreSql.AdoNet/Connections/TinyPostgreSqlAdoNetWorkerConnectionFactory.cs +++ b/src/TinyEvents.PostgreSql.AdoNet/Connections/TinyPostgreSqlAdoNetWorkerConnectionFactory.cs @@ -22,6 +22,8 @@ public TinyPostgreSqlAdoNetWorkerConnectionFactory( throw new ArgumentNullException(nameof(serviceProvider)); } + options.ValidateWorkerConfiguration(); + this.options = options; this.serviceProvider = serviceProvider; } diff --git a/src/TinyEvents.PostgreSql.AdoNet/TinyEventsPostgreSqlAdoNetOptions.cs b/src/TinyEvents.PostgreSql.AdoNet/TinyEventsPostgreSqlAdoNetOptions.cs index d6bf2b0..a79b5e5 100644 --- a/src/TinyEvents.PostgreSql.AdoNet/TinyEventsPostgreSqlAdoNetOptions.cs +++ b/src/TinyEvents.PostgreSql.AdoNet/TinyEventsPostgreSqlAdoNetOptions.cs @@ -42,12 +42,17 @@ internal ValueTask CreateWorkerConnectionAsync( throw new ArgumentNullException(nameof(serviceProvider)); } + ValidateWorkerConfiguration(); + + return workerConnectionFactory!(serviceProvider, cancellationToken); + } + + internal void ValidateWorkerConfiguration() + { if (workerConnectionFactory is null) { throw new InvalidOperationException( "An ADO.NET worker connection factory is required. Configure UseWorkerConnectionFactory(...) for outbox claiming and marking operations."); } - - return workerConnectionFactory(serviceProvider, cancellationToken); } } diff --git a/src/TinyEvents.PostgreSql.EntityFrameworkCore/TinyPostgreSqlEfCoreOutboxStore.cs b/src/TinyEvents.PostgreSql.EntityFrameworkCore/TinyPostgreSqlEfCoreOutboxStore.cs index 982b2d4..b1a0958 100644 --- a/src/TinyEvents.PostgreSql.EntityFrameworkCore/TinyPostgreSqlEfCoreOutboxStore.cs +++ b/src/TinyEvents.PostgreSql.EntityFrameworkCore/TinyPostgreSqlEfCoreOutboxStore.cs @@ -8,6 +8,8 @@ namespace TinyEvents.PostgreSql.EntityFrameworkCore; internal sealed class TinyPostgreSqlEfCoreOutboxStore : ITinyOutboxStore where TDbContext : DbContext { + private const string PostgreSqlProviderName = "Npgsql.EntityFrameworkCore.PostgreSQL"; + private readonly TDbContext dbContext; private readonly TinyPostgreSqlEfCoreTableName tableName; @@ -25,10 +27,24 @@ public TinyPostgreSqlEfCoreOutboxStore( throw new ArgumentNullException(nameof(options)); } + ValidateDatabaseProvider(dbContext); + this.dbContext = dbContext; tableName = TinyPostgreSqlEfCoreTableName.Parse(options.TableName); } + private static void ValidateDatabaseProvider(TDbContext dbContext) + { + if (!string.Equals( + dbContext.Database.ProviderName, + PostgreSqlProviderName, + StringComparison.Ordinal)) + { + throw new InvalidOperationException( + "The TinyEvents PostgreSQL EF Core provider requires the Npgsql EF Core provider. Configure the DbContext with UseNpgsql(...)."); + } + } + public async ValueTask> ClaimPendingAsync( int maxCount, string workerId, diff --git a/src/TinyEvents.SqlServer.AdoNet/Connections/TinySqlServerAdoNetWorkerConnectionFactory.cs b/src/TinyEvents.SqlServer.AdoNet/Connections/TinySqlServerAdoNetWorkerConnectionFactory.cs index cdc615b..d1fbd87 100644 --- a/src/TinyEvents.SqlServer.AdoNet/Connections/TinySqlServerAdoNetWorkerConnectionFactory.cs +++ b/src/TinyEvents.SqlServer.AdoNet/Connections/TinySqlServerAdoNetWorkerConnectionFactory.cs @@ -22,6 +22,8 @@ public TinySqlServerAdoNetWorkerConnectionFactory( throw new ArgumentNullException(nameof(serviceProvider)); } + options.ValidateWorkerConfiguration(); + this.options = options; this.serviceProvider = serviceProvider; } diff --git a/src/TinyEvents.SqlServer.AdoNet/TinyEventsSqlServerAdoNetOptions.cs b/src/TinyEvents.SqlServer.AdoNet/TinyEventsSqlServerAdoNetOptions.cs index 7adc4f1..2da2d25 100644 --- a/src/TinyEvents.SqlServer.AdoNet/TinyEventsSqlServerAdoNetOptions.cs +++ b/src/TinyEvents.SqlServer.AdoNet/TinyEventsSqlServerAdoNetOptions.cs @@ -42,12 +42,17 @@ internal ValueTask CreateWorkerConnectionAsync( throw new ArgumentNullException(nameof(serviceProvider)); } + ValidateWorkerConfiguration(); + + return workerConnectionFactory!(serviceProvider, cancellationToken); + } + + internal void ValidateWorkerConfiguration() + { if (workerConnectionFactory is null) { throw new InvalidOperationException( "An ADO.NET worker connection factory is required. Configure UseWorkerConnectionFactory(...) for outbox claiming and marking operations."); } - - return workerConnectionFactory(serviceProvider, cancellationToken); } } diff --git a/src/TinyEvents.SqlServer.EntityFrameworkCore/TinySqlServerEfCoreOutboxStore.cs b/src/TinyEvents.SqlServer.EntityFrameworkCore/TinySqlServerEfCoreOutboxStore.cs index 77ef5c7..5f59c41 100644 --- a/src/TinyEvents.SqlServer.EntityFrameworkCore/TinySqlServerEfCoreOutboxStore.cs +++ b/src/TinyEvents.SqlServer.EntityFrameworkCore/TinySqlServerEfCoreOutboxStore.cs @@ -8,6 +8,8 @@ namespace TinyEvents.SqlServer.EntityFrameworkCore; internal sealed class TinySqlServerEfCoreOutboxStore : ITinyOutboxStore where TDbContext : DbContext { + private const string SqlServerProviderName = "Microsoft.EntityFrameworkCore.SqlServer"; + private readonly TDbContext dbContext; private readonly TinySqlServerEfCoreTableName tableName; @@ -25,10 +27,24 @@ public TinySqlServerEfCoreOutboxStore( throw new ArgumentNullException(nameof(options)); } + ValidateDatabaseProvider(dbContext); + this.dbContext = dbContext; tableName = TinySqlServerEfCoreTableName.Parse(options.TableName); } + private static void ValidateDatabaseProvider(TDbContext dbContext) + { + if (!string.Equals( + dbContext.Database.ProviderName, + SqlServerProviderName, + StringComparison.Ordinal)) + { + throw new InvalidOperationException( + "The TinyEvents SQL Server EF Core provider requires the SQL Server EF Core provider. Configure the DbContext with UseSqlServer(...)."); + } + } + public async ValueTask> ClaimPendingAsync( int maxCount, string workerId, diff --git a/src/TinyEvents.Worker/DependencyInjection/TinyEventsWorkerServiceCollectionExtensions.cs b/src/TinyEvents.Worker/DependencyInjection/TinyEventsWorkerServiceCollectionExtensions.cs index d54f47f..6dcd2f5 100644 --- a/src/TinyEvents.Worker/DependencyInjection/TinyEventsWorkerServiceCollectionExtensions.cs +++ b/src/TinyEvents.Worker/DependencyInjection/TinyEventsWorkerServiceCollectionExtensions.cs @@ -15,14 +15,32 @@ public static IServiceCollection AddTinyEventsWorker( throw new ArgumentNullException(nameof(services)); } - var workerOptions = CreateOptions(configure); - services.TryAddSingleton(workerOptions); + var workerOptions = ConfigureOptions(services, configure); services.TryAddEnumerable(ServiceDescriptor.Singleton()); services.ConfigureTinyEventsForWorker(workerOptions); return services; } + private static TinyEventsWorkerOptions ConfigureOptions( + IServiceCollection services, + Action? configure) + { + var existingOptions = services + .LastOrDefault(descriptor => descriptor.ServiceType == typeof(TinyEventsWorkerOptions)) + ?.ImplementationInstance as TinyEventsWorkerOptions; + + if (existingOptions is not null) + { + configure?.Invoke(existingOptions); + return existingOptions; + } + + var options = CreateOptions(configure); + services.TryAddSingleton(options); + return options; + } + private static TinyEventsWorkerOptions CreateOptions(Action? configure) { var options = new TinyEventsWorkerOptions(); diff --git a/src/TinyEvents.Worker/Properties/AssemblyInfo.cs b/src/TinyEvents.Worker/Properties/AssemblyInfo.cs new file mode 100644 index 0000000..f74e701 --- /dev/null +++ b/src/TinyEvents.Worker/Properties/AssemblyInfo.cs @@ -0,0 +1,3 @@ +using System.Runtime.CompilerServices; + +[assembly: InternalsVisibleTo("TinyEvents.Worker.Tests")] diff --git a/src/TinyEvents.Worker/Startup/TinyEventsWorkerStartupValidator.cs b/src/TinyEvents.Worker/Startup/TinyEventsWorkerStartupValidator.cs new file mode 100644 index 0000000..73a51ff --- /dev/null +++ b/src/TinyEvents.Worker/Startup/TinyEventsWorkerStartupValidator.cs @@ -0,0 +1,21 @@ +using Microsoft.Extensions.DependencyInjection; + +namespace TinyEvents.Worker; + +internal sealed class TinyEventsWorkerStartupValidator +{ + private readonly IServiceScopeFactory scopeFactory; + + public TinyEventsWorkerStartupValidator(IServiceScopeFactory scopeFactory) + { + this.scopeFactory = scopeFactory + ?? throw new ArgumentNullException(nameof(scopeFactory)); + } + + public void ValidateConfiguration() + { + using var scope = scopeFactory.CreateScope(); + + _ = scope.ServiceProvider.GetRequiredService(); + } +} diff --git a/src/TinyEvents.Worker/TinyEventsBackgroundService.cs b/src/TinyEvents.Worker/TinyEventsBackgroundService.cs index a0a9206..52da3ff 100644 --- a/src/TinyEvents.Worker/TinyEventsBackgroundService.cs +++ b/src/TinyEvents.Worker/TinyEventsBackgroundService.cs @@ -8,6 +8,7 @@ namespace TinyEvents.Worker; public sealed class TinyEventsBackgroundService : BackgroundService { private readonly IServiceScopeFactory scopeFactory; + private readonly TinyEventsWorkerStartupValidator startupValidator; private readonly TinyEventsWorkerOptions options; private readonly ILogger logger; @@ -39,6 +40,7 @@ public TinyEventsBackgroundService( } this.scopeFactory = scopeFactory; + startupValidator = new TinyEventsWorkerStartupValidator(scopeFactory); this.options = options; this.logger = logger; } @@ -53,6 +55,10 @@ public async ValueTask ProcessOnceAsync(CancellationToken cancellationToken = de protected override async Task ExecuteAsync(CancellationToken stoppingToken) { + stoppingToken.ThrowIfCancellationRequested(); + + startupValidator.ValidateConfiguration(); + var failures = new TinyEventsWorkerFailureTracker(); while (!stoppingToken.IsCancellationRequested) diff --git a/tests/TinyEvents.PostgreSql.AdoNet.Tests/TinyPostgreSqlAdoNetOptionsTests.cs b/tests/TinyEvents.PostgreSql.AdoNet.Tests/TinyPostgreSqlAdoNetOptionsTests.cs index 6f2d173..001d50d 100644 --- a/tests/TinyEvents.PostgreSql.AdoNet.Tests/TinyPostgreSqlAdoNetOptionsTests.cs +++ b/tests/TinyEvents.PostgreSql.AdoNet.Tests/TinyPostgreSqlAdoNetOptionsTests.cs @@ -90,12 +90,10 @@ public async Task Worker_connection_factory_does_not_open_already_open_connectio } [Fact] - public async Task Worker_connection_factory_fails_clearly_when_delegate_is_missing() + public void Worker_connection_factory_rejects_missing_delegate() { - var factory = NewFactory(new TinyEventsPostgreSqlAdoNetOptions()); - - var exception = await Assert.ThrowsAsync( - async () => await factory.CreateOpenConnectionAsync(CancellationToken.None)); + var exception = Assert.Throws( + () => NewFactory(new TinyEventsPostgreSqlAdoNetOptions())); Assert.Contains("Configure UseWorkerConnectionFactory(...)", exception.Message); } diff --git a/tests/TinyEvents.PostgreSql.AdoNet.Tests/TinyPostgreSqlAdoNetWriterTests.cs b/tests/TinyEvents.PostgreSql.AdoNet.Tests/TinyPostgreSqlAdoNetWriterTests.cs index 184b254..62649b6 100644 --- a/tests/TinyEvents.PostgreSql.AdoNet.Tests/TinyPostgreSqlAdoNetWriterTests.cs +++ b/tests/TinyEvents.PostgreSql.AdoNet.Tests/TinyPostgreSqlAdoNetWriterTests.cs @@ -42,6 +42,41 @@ public void Use_postgre_sql_ado_net_outbox_registers_provider_services() provider.GetRequiredService()); } + [Fact] + public void Resolving_processor_rejects_missing_worker_connection_factory() + { + var services = new ServiceCollection(); + services.UsePostgreSqlAdoNetOutbox(_ => { }); + using var provider = services.BuildServiceProvider(); + using var scope = provider.CreateScope(); + + var exception = Assert.Throws( + () => scope.ServiceProvider.GetRequiredService()); + + Assert.Contains("Configure UseWorkerConnectionFactory(...)", exception.Message); + } + + [Fact] + public void Resolving_processor_does_not_invoke_worker_connection_factory() + { + var factoryCalls = 0; + var services = new ServiceCollection(); + services.UsePostgreSqlAdoNetOutbox(options => + { + options.UseWorkerConnectionFactory((_, _) => + { + factoryCalls++; + return new ValueTask(new RecordingConnection()); + }); + }); + using var provider = services.BuildServiceProvider(); + using var scope = provider.CreateScope(); + + _ = scope.ServiceProvider.GetRequiredService(); + + Assert.Equal(0, factoryCalls); + } + [Fact] public void Use_postgre_sql_ado_net_outbox_does_not_register_unit_of_work() { diff --git a/tests/TinyEvents.PostgreSql.EntityFrameworkCore.Tests/TinyPostgreSqlEfCoreOutboxStoreTests.cs b/tests/TinyEvents.PostgreSql.EntityFrameworkCore.Tests/TinyPostgreSqlEfCoreOutboxStoreTests.cs index e29070d..360e885 100644 --- a/tests/TinyEvents.PostgreSql.EntityFrameworkCore.Tests/TinyPostgreSqlEfCoreOutboxStoreTests.cs +++ b/tests/TinyEvents.PostgreSql.EntityFrameworkCore.Tests/TinyPostgreSqlEfCoreOutboxStoreTests.cs @@ -41,6 +41,35 @@ public void Use_postgre_sql_entity_framework_core_outbox_registers_configured_op Assert.Equal("app.MyOutbox", options.TableName); } + [Fact] + public void Resolving_processor_rejects_non_postgre_sql_db_context() + { + var services = new ServiceCollection(); + services.AddDbContext(options => + options.UseInMemoryDatabase(Guid.NewGuid().ToString())); + services.UsePostgreSqlEntityFrameworkCoreOutbox(); + using var provider = services.BuildServiceProvider(); + using var scope = provider.CreateScope(); + + var exception = Assert.Throws( + () => scope.ServiceProvider.GetRequiredService()); + + Assert.Contains("requires the Npgsql EF Core provider", exception.Message); + } + + [Fact] + public void Resolving_processor_accepts_postgre_sql_db_context_without_connecting() + { + var services = new ServiceCollection(); + services.AddDbContext(options => + options.UseNpgsql("Host=localhost;Database=TinyEvents")); + services.UsePostgreSqlEntityFrameworkCoreOutbox(); + using var provider = services.BuildServiceProvider(); + using var scope = provider.CreateScope(); + + _ = scope.ServiceProvider.GetRequiredService(); + } + [Fact] public void EF_store_does_not_implement_writer() { @@ -194,7 +223,7 @@ private static void AssertService(IServiceCollection private static TestDbContext NewTestDbContext() { var options = new DbContextOptionsBuilder() - .UseInMemoryDatabase(Guid.NewGuid().ToString()) + .UseNpgsql("Host=localhost;Database=TinyEvents") .Options; return new TestDbContext(options); diff --git a/tests/TinyEvents.SqlServer.AdoNet.Tests/TinySqlServerAdoNetProviderTests.cs b/tests/TinyEvents.SqlServer.AdoNet.Tests/TinySqlServerAdoNetProviderTests.cs index b85dc6f..2bd8219 100644 --- a/tests/TinyEvents.SqlServer.AdoNet.Tests/TinySqlServerAdoNetProviderTests.cs +++ b/tests/TinyEvents.SqlServer.AdoNet.Tests/TinySqlServerAdoNetProviderTests.cs @@ -42,6 +42,41 @@ public void UseSqlServerAdoNetOutbox_accepts_worker_connection_factory_delegate( && descriptor.ImplementationType == typeof(TinySqlServerAdoNetWorkerConnectionFactory)); } + [Fact] + public void Resolving_processor_rejects_missing_worker_connection_factory() + { + var services = new ServiceCollection(); + services.UseSqlServerAdoNetOutbox(_ => { }); + using var provider = services.BuildServiceProvider(); + using var scope = provider.CreateScope(); + + var exception = Assert.Throws( + () => scope.ServiceProvider.GetRequiredService()); + + Assert.Contains("Configure UseWorkerConnectionFactory(...)", exception.Message); + } + + [Fact] + public void Resolving_processor_does_not_invoke_worker_connection_factory() + { + var factoryCalls = 0; + var services = new ServiceCollection(); + services.UseSqlServerAdoNetOutbox(options => + { + options.UseWorkerConnectionFactory((_, _) => + { + factoryCalls++; + return new ValueTask(new RecordingConnection()); + }); + }); + using var provider = services.BuildServiceProvider(); + using var scope = provider.CreateScope(); + + _ = scope.ServiceProvider.GetRequiredService(); + + Assert.Equal(0, factoryCalls); + } + [Fact] public void UseSqlServerAdoNetOutbox_does_not_register_unit_of_work() { diff --git a/tests/TinyEvents.SqlServer.EntityFrameworkCore.Tests/TinyEvents.SqlServer.EntityFrameworkCore.Tests.csproj b/tests/TinyEvents.SqlServer.EntityFrameworkCore.Tests/TinyEvents.SqlServer.EntityFrameworkCore.Tests.csproj index dfe3567..8b27be0 100644 --- a/tests/TinyEvents.SqlServer.EntityFrameworkCore.Tests/TinyEvents.SqlServer.EntityFrameworkCore.Tests.csproj +++ b/tests/TinyEvents.SqlServer.EntityFrameworkCore.Tests/TinyEvents.SqlServer.EntityFrameworkCore.Tests.csproj @@ -8,9 +8,9 @@ + - diff --git a/tests/TinyEvents.SqlServer.EntityFrameworkCore.Tests/TinySqlServerEfCoreProviderTests.cs b/tests/TinyEvents.SqlServer.EntityFrameworkCore.Tests/TinySqlServerEfCoreProviderTests.cs index 2ac497e..e371bf9 100644 --- a/tests/TinyEvents.SqlServer.EntityFrameworkCore.Tests/TinySqlServerEfCoreProviderTests.cs +++ b/tests/TinyEvents.SqlServer.EntityFrameworkCore.Tests/TinySqlServerEfCoreProviderTests.cs @@ -127,6 +127,35 @@ public void Use_sql_server_entity_framework_core_outbox_registers_writer_and_sto && descriptor.ImplementationType == typeof(TinySqlServerEfCoreOutboxStore)); } + [Fact] + public void Resolving_processor_rejects_non_sql_server_db_context() + { + var services = new ServiceCollection(); + services.AddDbContext(options => + options.UseInMemoryDatabase(Guid.NewGuid().ToString())); + services.UseSqlServerEntityFrameworkCoreOutbox(); + using var provider = services.BuildServiceProvider(); + using var scope = provider.CreateScope(); + + var exception = Assert.Throws( + () => scope.ServiceProvider.GetRequiredService()); + + Assert.Contains("requires the SQL Server EF Core provider", exception.Message); + } + + [Fact] + public void Resolving_processor_accepts_sql_server_db_context_without_connecting() + { + var services = new ServiceCollection(); + services.AddDbContext(options => + options.UseSqlServer("Server=localhost;Database=TinyEvents;Integrated Security=true;TrustServerCertificate=true")); + services.UseSqlServerEntityFrameworkCoreOutbox(); + using var provider = services.BuildServiceProvider(); + using var scope = provider.CreateScope(); + + _ = scope.ServiceProvider.GetRequiredService(); + } + [Fact] public void EF_store_does_not_implement_writer() { diff --git a/tests/TinyEvents.Worker.Tests/TinyEventsWorkerStartupValidatorTests.cs b/tests/TinyEvents.Worker.Tests/TinyEventsWorkerStartupValidatorTests.cs new file mode 100644 index 0000000..e1e72d8 --- /dev/null +++ b/tests/TinyEvents.Worker.Tests/TinyEventsWorkerStartupValidatorTests.cs @@ -0,0 +1,197 @@ +using Microsoft.Extensions.DependencyInjection; +using TinyEvents.Worker; +using Xunit; + +namespace TinyEvents.Worker.Tests; + +public sealed class TinyEventsWorkerStartupValidatorTests +{ + [Fact] + public void Constructor_rejects_null_scope_factory() + { + Assert.Throws( + () => new TinyEventsWorkerStartupValidator(null!)); + } + + [Fact] + public void Validate_configuration_resolves_processor_without_invoking_it() + { + RecordingProcessor.Reset(); + var services = new ServiceCollection(); + services.AddScoped(); + using var provider = services.BuildServiceProvider(); + var scopeFactory = provider.GetRequiredService(); + var validator = new TinyEventsWorkerStartupValidator(scopeFactory); + + validator.ValidateConfiguration(); + + Assert.Equal(1, RecordingProcessor.CreatedCount); + Assert.Equal(0, RecordingProcessor.CallCount); + Assert.Equal(1, RecordingProcessor.DisposedCount); + } + + [Fact] + public void Validate_configuration_rejects_missing_processor() + { + var services = new ServiceCollection(); + using var provider = services.BuildServiceProvider(); + var validator = CreateValidator(provider); + + var exception = Assert.Throws( + validator.ValidateConfiguration); + + Assert.Contains(nameof(ITinyOutboxProcessor), exception.Message, StringComparison.Ordinal); + } + + [Fact] + public void Validate_configuration_rejects_missing_store() + { + var services = new ServiceCollection(); + services.UseTinyEvents(); + using var provider = services.BuildServiceProvider(); + var validator = CreateValidator(provider); + + var exception = Assert.Throws( + validator.ValidateConfiguration); + + Assert.Contains( + typeof(TinyOutboxProcessor).FullName!, + exception.Message, + StringComparison.Ordinal); + } + + [Fact] + public void Validate_configuration_rejects_conflicting_dispatchers() + { + var services = CreateProcessorServices(); + services.AddSingleton( + new TinyEventDispatcher("Shared.Event")); + services.AddSingleton( + new TinyEventDispatcher("Shared.Event")); + using var provider = services.BuildServiceProvider(); + var validator = CreateValidator(provider); + + var exception = Assert.Throws( + validator.ValidateConfiguration); + + Assert.Contains("Shared.Event", exception.Message, StringComparison.Ordinal); + Assert.Contains(typeof(FirstEvent).FullName!, exception.Message, StringComparison.Ordinal); + Assert.Contains(typeof(SecondEvent).FullName!, exception.Message, StringComparison.Ordinal); + } + + [Fact] + public void Validate_configuration_does_not_resolve_or_invoke_consumers() + { + var consumerCreatedCount = 0; + var services = CreateProcessorServices(); + services.AddScoped>(_ => + { + consumerCreatedCount++; + return new NoOpConsumer(); + }); + services.AddSingleton( + new TinyEventDispatcher(typeof(FirstEvent).FullName!)); + using var provider = services.BuildServiceProvider(); + var validator = CreateValidator(provider); + + validator.ValidateConfiguration(); + + Assert.Equal(0, consumerCreatedCount); + } + + private static ServiceCollection CreateProcessorServices() + { + var services = new ServiceCollection(); + services.UseTinyEvents(); + services.AddSingleton(); + return services; + } + + private static TinyEventsWorkerStartupValidator CreateValidator( + IServiceProvider serviceProvider) + { + var scopeFactory = serviceProvider.GetRequiredService(); + return new TinyEventsWorkerStartupValidator(scopeFactory); + } + + private sealed class RecordingProcessor : ITinyOutboxProcessor, IDisposable + { + public RecordingProcessor() + { + CreatedCount++; + } + + public static int CreatedCount { get; private set; } + + public static int CallCount { get; private set; } + + public static int DisposedCount { get; private set; } + + public static void Reset() + { + CreatedCount = 0; + CallCount = 0; + DisposedCount = 0; + } + + public ValueTask ProcessPendingAsync(CancellationToken cancellationToken = default) + { + CallCount++; + return ValueTask.CompletedTask; + } + + public void Dispose() + { + DisposedCount++; + } + } + + private sealed record FirstEvent; + + private sealed record SecondEvent; + + private sealed class NoOpConsumer : IEventConsumer + { + public ValueTask ConsumeAsync( + FirstEvent @event, + CancellationToken cancellationToken) + { + return ValueTask.CompletedTask; + } + } + + private sealed class NoOpStore : ITinyOutboxStore + { + public ValueTask> ClaimPendingAsync( + int maxCount, + string workerId, + DateTimeOffset now, + TimeSpan claimTimeout, + CancellationToken cancellationToken) + { + return ValueTask.FromResult>( + Array.Empty()); + } + + public ValueTask MarkProcessedAsync( + Guid messageId, + string workerId, + DateTimeOffset processedAtUtc, + CancellationToken cancellationToken) + { + return ValueTask.CompletedTask; + } + + public ValueTask MarkFailedAsync( + Guid messageId, + string workerId, + string error, + int attemptCount, + DateTimeOffset? nextAttemptAtUtc, + CancellationToken cancellationToken) + { + return ValueTask.CompletedTask; + } + } + +} diff --git a/tests/TinyEvents.Worker.Tests/TinyEventsWorkerTests.cs b/tests/TinyEvents.Worker.Tests/TinyEventsWorkerTests.cs index cad558a..1505ee3 100644 --- a/tests/TinyEvents.Worker.Tests/TinyEventsWorkerTests.cs +++ b/tests/TinyEvents.Worker.Tests/TinyEventsWorkerTests.cs @@ -130,6 +130,34 @@ public void Add_tiny_events_worker_preserves_core_retry_configuration() Assert.Equal(TimeSpan.FromSeconds(17), coreOptions.RetryDelay); } + [Fact] + public void Repeated_worker_registration_keeps_worker_and_core_options_consistent() + { + var services = new ServiceCollection(); + services.AddTinyEventsWorker(options => + { + options.WorkerId = "first-worker"; + options.BatchSize = 3; + options.PollingInterval = TimeSpan.FromSeconds(1); + }); + services.AddTinyEventsWorker(options => + { + options.WorkerId = "second-worker"; + options.BatchSize = 7; + options.PollingInterval = TimeSpan.FromSeconds(2); + }); + + using var provider = services.BuildServiceProvider(); + var workerOptions = provider.GetRequiredService(); + var coreOptions = provider.GetRequiredService(); + + Assert.Equal("second-worker", workerOptions.WorkerId); + Assert.Equal(workerOptions.WorkerId, coreOptions.WorkerId); + Assert.Equal(7, workerOptions.BatchSize); + Assert.Equal(workerOptions.BatchSize, coreOptions.BatchSize); + Assert.Equal(TimeSpan.FromSeconds(2), workerOptions.PollingInterval); + } + [Fact] public void Add_tiny_events_worker_rejects_empty_configured_worker_id() { @@ -212,6 +240,43 @@ public async Task Background_service_process_once_resolves_processor_from_scope( Assert.Equal(1, RecordingProcessor.CallCount); } + [Fact] + public async Task Background_service_fails_startup_before_polling_when_processor_is_missing() + { + var logger = new RecordingLogger(); + var services = new ServiceCollection(); + services.AddSingleton(new TinyEventsWorkerOptions()); + services.AddSingleton>(logger); + services.AddSingleton(); + using var provider = services.BuildServiceProvider(); + var worker = provider.GetRequiredService(); + + var exception = await Assert.ThrowsAsync( + () => worker.StartAsync(CancellationToken.None)); + + Assert.Contains(nameof(ITinyOutboxProcessor), exception.Message, StringComparison.Ordinal); + Assert.Empty(logger.Entries); + } + + [Fact] + public async Task Background_service_skips_startup_validation_when_already_stopping() + { + StartupTrackingProcessor.ConstructionCount = 0; + var services = new ServiceCollection(); + services.AddScoped(); + services.AddSingleton(new TinyEventsWorkerOptions()); + services.AddSingleton(); + using var provider = services.BuildServiceProvider(); + var worker = provider.GetRequiredService(); + using var cancellation = new CancellationTokenSource(); + cancellation.Cancel(); + + await Assert.ThrowsAnyAsync( + () => worker.StartAsync(cancellation.Token)); + + Assert.Equal(0, StartupTrackingProcessor.ConstructionCount); + } + [Fact] public async Task Background_service_creates_scope_per_processing_iteration() { @@ -356,6 +421,21 @@ public ValueTask ProcessPendingAsync(CancellationToken cancellationToken = defau } } + private sealed class StartupTrackingProcessor : ITinyOutboxProcessor + { + public StartupTrackingProcessor() + { + ConstructionCount++; + } + + public static int ConstructionCount { get; set; } + + public ValueTask ProcessPendingAsync(CancellationToken cancellationToken = default) + { + return ValueTask.CompletedTask; + } + } + private sealed class ScopedProcessor : ITinyOutboxProcessor { private readonly Guid instanceId = Guid.NewGuid();