Philiprehberger.EventBus 0.6.0

dotnet add package Philiprehberger.EventBus --version 0.6.0
                    
NuGet\Install-Package Philiprehberger.EventBus -Version 0.6.0
                    
This command is intended to be used within the Package Manager Console in Visual Studio, as it uses the NuGet module's version of Install-Package.
<PackageReference Include="Philiprehberger.EventBus" Version="0.6.0" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="Philiprehberger.EventBus" Version="0.6.0" />
                    
Directory.Packages.props
<PackageReference Include="Philiprehberger.EventBus" />
                    
Project file
For projects that support Central Package Management (CPM), copy this XML node into the solution Directory.Packages.props file to version the package.
paket add Philiprehberger.EventBus --version 0.6.0
                    
#r "nuget: Philiprehberger.EventBus, 0.6.0"
                    
#r directive can be used in F# Interactive and Polyglot Notebooks. Copy this into the interactive tool or source code of the script to reference the package.
#:package Philiprehberger.EventBus@0.6.0
                    
#:package directive can be used in C# file-based apps starting in .NET 10 preview 4. Copy this into a .cs file before any lines of code to reference the package.
#addin nuget:?package=Philiprehberger.EventBus&version=0.6.0
                    
Install as a Cake Addin
#tool nuget:?package=Philiprehberger.EventBus&version=0.6.0
                    
Install as a Cake Tool

Philiprehberger.EventBus

CI NuGet Last updated

Philiprehberger.EventBus

In-process publish/subscribe event bus with middleware pipeline, dead-letter queue, event replay, and Microsoft DI integration.

Installation

dotnet add package Philiprehberger.EventBus

Usage

using Philiprehberger.EventBus;

var bus = new EventBus();

using var subscription = bus.Subscribe<OrderPlaced>(async (e, ct) =>
{
    Console.WriteLine($"Order {e.OrderId} placed");
});

await bus.PublishAsync(new OrderPlaced(OrderId: 42));

record OrderPlaced(int OrderId);

Publish and Subscribe

using Philiprehberger.EventBus;

var bus = new EventBus();

// Subscribe returns an IDisposable — dispose it to unsubscribe
using var sub = bus.Subscribe<UserRegistered>(async (e, ct) =>
{
    await SendWelcomeEmailAsync(e.Email, ct);
});

// Publish fires all handlers concurrently
await bus.PublishAsync(new UserRegistered("user@example.com"));

record UserRegistered(string Email);

Handler Priority

using Philiprehberger.EventBus;

var bus = new EventBus();

// Lower priority number executes first
bus.Subscribe<OrderPlaced>((e, ct) =>
{
    Console.WriteLine("Validate order");
    return Task.CompletedTask;
}, priority: 10);

bus.Subscribe<OrderPlaced>((e, ct) =>
{
    Console.WriteLine("Send confirmation email");
    return Task.CompletedTask;
}, priority: 20);

await bus.PublishAsync(new OrderPlaced(1));
// Output: Validate order, then Send confirmation email

record OrderPlaced(int OrderId);

Handler Filtering

using Philiprehberger.EventBus;

var bus = new EventBus();

// Only handle high-value orders
bus.Subscribe<OrderPlaced>(
    (e, ct) =>
    {
        Console.WriteLine($"High-value order: {e.OrderId}");
        return Task.CompletedTask;
    },
    filter: e => e.Total > 1000);

await bus.PublishAsync(new OrderPlaced(1, 500));   // Skipped
await bus.PublishAsync(new OrderPlaced(2, 2000));  // Handled

record OrderPlaced(int OrderId, decimal Total);

Error Handling

using Philiprehberger.EventBus;

var bus = new EventBus(new EventBusOptions
{
    ThrowOnHandlerError = false,
    OnHandlerError = ex => Console.Error.WriteLine($"Handler failed: {ex.Message}")
});

bus.Subscribe<OrderPlaced>((_, _) => throw new InvalidOperationException("oops"));

await bus.PublishAsync(new OrderPlaced(1));
// Logs "Handler failed: oops" without propagating the exception

record OrderPlaced(int OrderId);

Handler Timeout

using Philiprehberger.EventBus;

var bus = new EventBus(new EventBusOptions
{
    ThrowOnHandlerError = true,
    HandlerTimeout = TimeSpan.FromSeconds(5)
});

bus.Subscribe<OrderPlaced>(async (e, ct) =>
{
    await ProcessOrderAsync(e.OrderId, ct);
});

// Throws TimeoutException if the handler exceeds 5 seconds
await bus.PublishAsync(new OrderPlaced(1));

record OrderPlaced(int OrderId);

Dead-Letter Queue

using Philiprehberger.EventBus;

var bus = new EventBus(new EventBusOptions
{
    ThrowOnHandlerError = false,
    OnDeadLetter = (evt, ex) =>
        Console.Error.WriteLine($"Dead letter: {evt.GetType().Name} failed with {ex.Message}")
});

bus.Subscribe<OrderPlaced>((_, _) => throw new InvalidOperationException("payment failed"));

await bus.PublishAsync(new OrderPlaced(1));
// Logs "Dead letter: OrderPlaced failed with payment failed"

record OrderPlaced(int OrderId);

Event Replay

using Philiprehberger.EventBus;

var bus = new EventBus();
bus.EnableHistory(maxEvents: 100);

bus.Subscribe<OrderPlaced>((e, _) =>
{
    Console.WriteLine($"Order {e.OrderId}");
    return Task.CompletedTask;
});

await bus.PublishAsync(new OrderPlaced(1));
await bus.PublishAsync(new OrderPlaced(2));
await bus.PublishAsync(new OrderPlaced(3));

// Re-publishes the 2 most recent events (OrderPlaced 2, then 3)
await bus.ReplayLastAsync(2);

record OrderPlaced(int OrderId);

Middleware

using Philiprehberger.EventBus;

var bus = new EventBus();

// Add logging middleware
bus.Use(async (context, next) =>
{
    Console.WriteLine($"Before: {context.EventType.Name}");
    await next();
    Console.WriteLine($"After: {context.EventType.Name}");
});

// Add timing middleware
bus.Use(async (context, next) =>
{
    var sw = System.Diagnostics.Stopwatch.StartNew();
    await next();
    sw.Stop();
    Console.WriteLine($"Handler took {sw.ElapsedMilliseconds}ms");
});

bus.Subscribe<OrderPlaced>((e, _) =>
{
    Console.WriteLine($"Processing order {e.OrderId}");
    return Task.CompletedTask;
});

await bus.PublishAsync(new OrderPlaced(1));

record OrderPlaced(int OrderId);

One-Time Subscription

using Philiprehberger.EventBus;

var bus = new EventBus();

// Handler is automatically unsubscribed after the first matching event
bus.SubscribeOnce<OrderPlaced>((e, _) =>
{
    Console.WriteLine($"First order placed: {e.OrderId}");
    return Task.CompletedTask;
});

await bus.PublishAsync(new OrderPlaced(1)); // Handled
await bus.PublishAsync(new OrderPlaced(2)); // Ignored — already unsubscribed

record OrderPlaced(int OrderId);

Await Next Event

using Philiprehberger.EventBus;

var bus = new EventBus();

// Start waiting before the event is published
var waitTask = bus.WaitForAsync<OrderPlaced>(
    filter: e => e.Total > 1000);

// Publish events from another part of the application
await bus.PublishAsync(new OrderPlaced(1, 500));   // Skipped by filter
await bus.PublishAsync(new OrderPlaced(2, 2000));  // Completes the wait

var order = await waitTask;
Console.WriteLine($"High-value order: {order.OrderId}");

record OrderPlaced(int OrderId, decimal Total);

Check Subscribers

using Philiprehberger.EventBus;

var bus = new EventBus();

Console.WriteLine(bus.HasSubscribers<OrderPlaced>()); // False
Console.WriteLine(bus.GetSubscriberCount<OrderPlaced>()); // 0

using var sub = bus.Subscribe<OrderPlaced>((_, _) => Task.CompletedTask);
Console.WriteLine(bus.HasSubscribers<OrderPlaced>()); // True
Console.WriteLine(bus.GetSubscriberCount<OrderPlaced>()); // 1

record OrderPlaced(int OrderId);

Bulk Unsubscribe

using Philiprehberger.EventBus;

var bus = new EventBus();

bus.Subscribe<OrderPlaced>((e, _) => { Console.WriteLine(e.OrderId); return Task.CompletedTask; });
bus.Subscribe<OrderPlaced>((e, _) => { Console.WriteLine("Audit: " + e.OrderId); return Task.CompletedTask; });

// Remove all handlers for a specific event type
bus.UnsubscribeAll<OrderPlaced>();

// Or remove all handlers for all event types
bus.UnsubscribeAll();

record OrderPlaced(int OrderId);

Inspect Event History

using Philiprehberger.EventBus;

var bus = new EventBus();
bus.EnableHistory(100);

await bus.PublishAsync(new OrderPlaced(1));
await bus.PublishAsync(new OrderPlaced(2));
await bus.PublishAsync(new OrderPlaced(3));

// Read recorded events without re-triggering handlers
IReadOnlyList<object> history = bus.GetHistory();
foreach (var evt in history)
{
    Console.WriteLine(((OrderPlaced)evt).OrderId);
}
// Output: 1, 2, 3

// Typed history filtered to a specific event type
IReadOnlyList<OrderPlaced> orders = bus.GetHistory<OrderPlaced>();

record OrderPlaced(int OrderId);

Toggle History

using Philiprehberger.EventBus;

var bus = new EventBus();

if (!bus.IsHistoryEnabled)
{
    bus.EnableHistory(100);
}

// ...later, when telemetry is no longer needed...
bus.DisableHistory();

DI Registration

using Philiprehberger.EventBus;

var builder = WebApplication.CreateBuilder(args);

builder.Services.AddEventBus(options =>
{
    options.ThrowOnHandlerError = true;
    options.MaxConcurrency = 4;
    options.HandlerTimeout = TimeSpan.FromSeconds(10);
    options.OnHandlerError = ex => Console.Error.WriteLine(ex);
});

var app = builder.Build();

Handler Classes

using Philiprehberger.EventBus;

public record OrderShipped(int OrderId, string TrackingNumber);

public class OrderShippedHandler : IEventHandler<OrderShipped>
{
    public async Task HandleAsync(OrderShipped @event, CancellationToken ct)
    {
        await NotifyCustomerAsync(@event.OrderId, @event.TrackingNumber, ct);
    }
}

// Subscribe an IEventHandler<T> instance directly (no DI required)
var bus = new EventBus();
var handler = new OrderShippedHandler();
using var sub = bus.Subscribe<OrderShipped>(handler);

API

IEventBus

Method Description
PublishAsync<T>(@event, ct) Publishes an event to all registered handlers for the type
Subscribe<T>(handler, priority, filter) Subscribes a handler function; returns IDisposable to unsubscribe
Subscribe<T>(IEventHandler<T>, priority, filter) Subscribes an IEventHandler<T> instance; returns IDisposable to unsubscribe
Use(middleware) Registers a middleware function that wraps every handler invocation
EnableHistory(maxEvents) Enables circular buffer event history with the specified capacity
DisableHistory() Disables history tracking and releases the buffer
IsHistoryEnabled bool property; true when history tracking is currently enabled
ReplayLastAsync(count, ct) Re-publishes the N most recent events from the history buffer
ClearHistory() Clears all events from the history buffer without disabling tracking
HasSubscribers<T>() Returns true if any handlers are registered for the event type
GetSubscriberCount<T>() Returns the number of handlers registered for the event type
UnsubscribeAll<T>() Removes all handlers for a specific event type
UnsubscribeAll() Removes all handlers for all event types
GetHistory() Returns a read-only snapshot of recorded events in chronological order
GetHistory<T>() Returns recorded events filtered to type T, oldest first
SubscribeOnce<T>(handler, filter) Subscribes a handler that auto-unsubscribes after one invocation
WaitForAsync<T>(filter, ct) Returns a Task<T> that completes with the next matching event

IEventHandler<T>

Method Description
HandleAsync(@event, ct) Handles an event of type T asynchronously

EventBusOptions

Property Type Default Description
ThrowOnHandlerError bool false Propagate handler exceptions to the publisher
MaxConcurrency int 0 Max concurrent handler invocations (0 = unlimited)
OnHandlerError Action<Exception>? null Callback invoked when any handler throws an exception
HandlerTimeout TimeSpan? null Timeout per handler invocation; throws TimeoutException if exceeded
OnDeadLetter Action<object, Exception>? null Callback invoked with the failed event and exception when a handler throws and ThrowOnHandlerError is false

Subscribe<T> Parameters

Parameter Type Default Description
handler Func<T, CancellationToken, Task> required The handler function to invoke
priority int 0 Execution priority; lower values execute first
filter Func<T, bool>? null Predicate evaluated before invoking; handler is skipped if it returns false

EventContext

Property Type Description
Event object The event instance being published
EventType Type The CLR type of the event
CancellationToken CancellationToken The cancellation token for the current publish operation
Items IDictionary<string, object> Dictionary for middleware to pass data along the pipeline

ServiceCollectionExtensions

Method Description
AddEventBus(configure?) Registers IEventBus as singleton and scans for IEventHandler<T> implementations

Development

dotnet build src/Philiprehberger.EventBus.csproj --configuration Release

Support

If you find this project useful:

Star the repo

🐛 Report issues

💡 Suggest features

❤️ Sponsor development

🌐 All Open Source Projects

💻 GitHub Profile

🔗 LinkedIn Profile

License

MIT

Product Compatible and additional computed target framework versions.
.NET net8.0 is compatible.  net8.0-android was computed.  net8.0-browser was computed.  net8.0-ios was computed.  net8.0-maccatalyst was computed.  net8.0-macos was computed.  net8.0-tvos was computed.  net8.0-windows was computed.  net9.0 was computed.  net9.0-android was computed.  net9.0-browser was computed.  net9.0-ios was computed.  net9.0-maccatalyst was computed.  net9.0-macos was computed.  net9.0-tvos was computed.  net9.0-windows was computed.  net10.0 was computed.  net10.0-android was computed.  net10.0-browser was computed.  net10.0-ios was computed.  net10.0-maccatalyst was computed.  net10.0-macos was computed.  net10.0-tvos was computed.  net10.0-windows was computed. 
Compatible target framework(s)
Included target framework(s) (in package)
Learn more about Target Frameworks and .NET Standard.

NuGet packages

This package is not used by any NuGet packages.

GitHub repositories

This package is not used by any popular GitHub repositories.

Version Downloads Last Updated
0.6.0 102 6/15/2026
0.5.0 114 4/13/2026
0.4.0 106 4/12/2026
0.3.0 116 4/1/2026
0.2.0 171 3/28/2026
0.1.3 110 3/27/2026
0.1.2 105 3/25/2026
0.1.1 110 3/23/2026
0.1.0 117 3/21/2026