Back to "Att bygga ett återanvändbart efemärt avrättningsbibliotek"

This is a viewer only at the moment see the article on how this works.

To update the preview hit Ctrl-Alt-R (or ⌘-Alt-R on Mac) or Enter to refresh. The Save icon lets you save the markdown file to disk

This is a preview from the server running through my markdig pipeline

Architecture ASP.NET Async DI Systems Design

Att bygga ett återanvändbart efemärt avrättningsbibliotek

Friday, 12 December 2025

Till Del 1: Eld och låt bli Ganska Glöm, Vi utforskade teorin bakom efemeral exekvering - avgränsade, privata, debuggable async arbetsflöden som minns precis tillräckligt för att vara användbar och sedan avdunsta.

Denna artikel förvandlar det mönstret till ett återanvändbart bibliotek som du kan släppa in i vilket .NET-projekt som helst.

NUGET!!!

Detta är nu i de flesta lucid.ephemerals Nuget paket också mer än 20 mestadels lucid.ephemerals mönster och "atomer".

Hämta Licens

Källfiler

Biblioteket är uppdelat i väl tilltagna filer:

Filens syfte |------|---------| | Efemeralalternativ.cs Konfiguration (konkurrent, fönsterstorlek, livslängd, signaler) | EfemeralOperation.cs till Intern drift spårning med signalstöd på | Snapshots.cs "Oförgängliga ögonblicksbilder som exponeras för konsumenter" | Signaler.cs Signalhändelser, spridning, begränsningar och den globala signalsänkan | EfemeralIdGenerator.cs på snabb XxHash64-baserad ID-generering | KoncurrencyGates.cs på fast och justerbar konvergensbegränsning | StringPatternMatcher.cs på ett mönster som passar för signalfiltrering | ParallellEfemeral.cs på grund av att det inte finns några statiska utvidgningsmetoder (EphemeralForEachAsync) | | Efemärarbetssamordnare.cs Långlivad arbetskösamordnare | EfemeralKeyedarbetssamordnare.cs till Per-nyckel sekventiellt utförande med rättvis schemaläggning | EfemeralResultatsamordnare.cs till resultat-kaptursamordnare variant på | Signalsändare.cs till Async-signal routing med mönstermatchning | BeroendeInjicering.cs Utbyggnadsmetoder för DI och fabriksimplementeringar | Exempel/SignalingHttpClient.cs på prov finkornig signalemission för HTTP-samtal

Och Övergripande tester som täcker alla kantfall.


Före och efter

Här är vad vi ersätter:

// ❌ Before: Fire-and-forget black hole
_ = Task.Run(() => ProcessAsync(item));
// No visibility. No debugging. No idea if it worked.

// ❌ Or: Blocking everything
await ProcessAsync(item);  // Hope you like waiting...

Och vad vi bygger:

// ✅ After: Trackable, bounded, debuggable
await coordinator.EnqueueAsync(item);

// Instant visibility
Console.WriteLine($"Pending: {coordinator.PendingCount}");
Console.WriteLine($"Active: {coordinator.ActiveCount}");
Console.WriteLine($"Failed: {coordinator.TotalFailed}");

// Full operation history
var snapshot = coordinator.GetSnapshot();
var failures = coordinator.GetFailed();

Samma async-utförande. Fullständig observerbarhet. Ingen användardata sparad.


Snabbstart

Det vanligaste mönstret - registrera en samordnare i DI och injicera den:

// Program.cs
services.AddEphemeralWorkCoordinator<TranslationRequest>(
    async (request, ct) => await TranslateAsync(request, ct),
    new EphemeralOptions { MaxConcurrency = 8 });

// Your service
public class TranslationService(EphemeralWorkCoordinator<TranslationRequest> coordinator)
{
    public async Task TranslateAsync(TranslationRequest request)
    {
        await coordinator.EnqueueAsync(request);
        // Returns immediately - work happens in background
    }

    public object GetStatus() => new
    {
        pending = coordinator.PendingCount,
        active = coordinator.ActiveCount,
        completed = coordinator.TotalCompleted,
        failed = coordinator.TotalFailed
    };
}

Vilken variant behöver jag?

┌─────────────────────────────────────────────────────────────────┐
│                    DECISION TREE                                │
├─────────────────────────────────────────────────────────────────┤
│                                                                 │
│  Processing a collection once?                                  │
│  └─► EphemeralForEachAsync<T> (ParallelEphemeral.cs)            │
│                                                                 │
│  Need a long-lived queue that accepts items over time?          │
│  └─► EphemeralWorkCoordinator<T>                                │
│                                                                 │
│  Need per-entity ordering (user commands, tenant jobs)?         │
│  └─► EphemeralKeyedWorkCoordinator<TKey, T>                     │
│                                                                 │
│  Need to capture results (fingerprints, summaries)?             │
│  └─► EphemeralResultCoordinator<TInput, TResult>                │
│                                                                 │
│  Need multiple coordinators with different configs?             │
│  └─► IEphemeralCoordinatorFactory<T> (like IHttpClientFactory)  │
│                                                                 │
│  Need dynamic concurrency adjustment at runtime?                │
│  └─► Set EnableDynamicConcurrency = true, call SetMaxConcurrency│
│                                                                 │
└─────────────────────────────────────────────────────────────────┘

Inställningsobjektet

Från Efemeralalternativ.cs:

public sealed class EphemeralOptions
{
    // Concurrency control
    public int MaxConcurrency { get; init; } = Environment.ProcessorCount;
    public int MaxConcurrencyPerKey { get; init; } = 1;
    public bool EnableDynamicConcurrency { get; init; } = false;

    // Window management
    public int MaxTrackedOperations { get; init; } = 200;
    public TimeSpan? MaxOperationLifetime { get; init; } = TimeSpan.FromMinutes(5);

    // Fair scheduling (keyed coordinator)
    public bool EnableFairScheduling { get; init; } = false;
    public int FairSchedulingThreshold { get; init; } = 10;

    // Signal-reactive processing
    public IReadOnlySet<string>? CancelOnSignals { get; init; }
    public IReadOnlySet<string>? DeferOnSignals { get; init; }
    public int MaxDeferAttempts { get; init; } = 10;
    public TimeSpan DeferCheckInterval { get; init; } = TimeSpan.FromMilliseconds(100);

    // Signal infrastructure
    public SignalSink? Signals { get; init; }
    public SignalConstraints? SignalConstraints { get; init; }
    public Action<SignalEvent>? OnSignal { get; init; }

    // Async signal handling
    public Func<SignalEvent, CancellationToken, Task>? OnSignalAsync { get; init; }
    public int MaxConcurrentSignalHandlers { get; init; } = 4;
    public int MaxQueuedSignals { get; init; } = 1000;

    // Observability
    public Action<IReadOnlyCollection<EphemeralOperationSnapshot>>? OnSample { get; init; }
}

Viktiga beslut om utformning

  • Maxlikviditet förvalt till CPU räkna - vettigt för CPU-bundet arbete. För I/O-bundet arbete, öka det.
  • AktiveraDynamicConcurrency möjliggör körtidsjustering via SetMaxConcurrency() - använder en egen grind istället för SemaphoreSlim.
  • AvbrytSignaler/avslutaSignaler@ info: whatsthis göra koordinatorer signal-reaktiva - de reagerar på omgivande system tillstånd (mönster matchning stöd */?/comma listor).
  • På signering är synkron; för async-utfläkt SignalDispatcher eller AsyncSignalProcessor Inuti handtaget.
  • Signalhalter förhindrar oändliga signalslingor med cykeldetektering och djupgränser.

Snapshot-skivorna

Från Snapshots.cs:

public sealed record EphemeralOperationSnapshot(
    long Id,
    DateTimeOffset Started,
    DateTimeOffset? Completed,
    string? Key,
    bool IsFaulted,
    Exception? Error,
    TimeSpan? Duration,
    IReadOnlyList<string>? Signals = null,
    bool IsPinned = false)
{
    public bool HasSignal(string signal) => Signals?.Contains(signal) == true;
}

// For result-capturing coordinators
public sealed record EphemeralOperationSnapshot<TResult>(
    long Id,
    DateTimeOffset Started,
    DateTimeOffset? Completed,
    string? Key,
    bool IsFaulted,
    Exception? Error,
    TimeSpan? Duration,
    TResult? Result,
    bool HasResult,
    IReadOnlyList<string>? Signals = null,
    bool IsPinned = false);

Det här är Endast metadataLägg märke till vad som är inte Här:

  • Ingen nyttolast
  • Inga indata
  • Inget användarinnehåll

Bara tillräckligt för att svara "vad hände, när, och fungerade det?" - inget mer.


Hur detta kan jämföras med andra metoder

.NET ger dig flera sätt att göra parallellarbete. Så här jämför Efemeralbiblioteket:

Parallellt.ForEachAsync (.NET 6+)

await Parallel.ForEachAsync(items,
    new ParallelOptions { MaxDegreeOfParallelism = 4 },
    async (item, ct) => await ProcessAsync(item, ct));

Bäst för: Enkel parallellbehandling av samlingar där du inte behöver synliggöras.

Vad den saknar:

  • Ingen operationsspårning
  • Ingen utförande per nyckel i följd
  • Ingen insyn i det som körs

Använd efemeral när: Du behöver felsökning/observerbarhet, per-nyckel beställning, eller signal-reactive processing.

TPL Dataflöde

var block = new ActionBlock<T>(
    async item => await ProcessAsync(item),
    new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = 4 });

foreach (var item in items)
    block.Post(item);

block.Complete();
await block.Completion;

Bäst för: Komplexa dataflöde pipelines med förgrening, sammanslagning, batching.

Vad den gör bra:

  • Rik rörledning sammansättning (länkblock tillsammans)
  • Inbyggd batching, transformering, sändning
  • Gränsad kapacitet med mottryck

Använd TPL Dataflöde när: Du behöver komplexa pipeline topologier (fan-out, fan-in, villkorlig routing).

Använd efemeral när: Du behöver operationsspårning, enklare API, eller signal-reaktiv samordning.

System. Threading.Channels

var channel = Channel.CreateBounded<T>(100);

// Producer
foreach (var item in items)
    await channel.Writer.WriteAsync(item);
channel.Writer.Complete();

// Consumer (multiple workers)
var workers = Enumerable.Range(0, 4).Select(async _ =>
{
    await foreach (var item in channel.Reader.ReadAllAsync())
        await ProcessAsync(item);
});
await Task.WhenAll(workers);

Bäst för: Producent-konsument mönster där du kontrollerar båda sidor.

Vad den gör bra:

  • Utmärkt prestanda
  • Mottryck via avgränsade kanaler
  • Separation av producenter och konsumenter

Använd kanaler när: Du bygger anpassad infrastruktur och behöver maximal kontroll.

Använd efemeral när: Du vill ha operationsspårning och observerbarhet utan pannplattan.

Polly Ordförande

var policy = Policy
    .Handle<HttpRequestException>()
    .WaitAndRetryAsync(3, attempt => TimeSpan.FromSeconds(Math.Pow(2, attempt)));

await policy.ExecuteAsync(() => ProcessAsync(item));

Bäst för: Resilienspolicyer (återgång, strömbrytare, timeout) för enskilda operationer.

Använd Polly när: Du behöver motståndskraft runt enskilda samtal.

Använd efemeral när: Du behöver samordning över många operationer med omgivningsmedvetenhet.

Kombinera dem: Använd Polly inuti din Ephemerala arbetskropp för per-operation resiliens.

MassTransit / NServiceBus

Bäst för: Distribuerade meddelanden över tjänster med varaktiga köer.

Använd meddelandebussar när: Arbetet måste överleva processen omstarter, spänna över flera tjänster, eller kräver garanterad leverans.

Använd efemeral när: Arbete är under bearbetning, behöver inte hållbarhet, och du vill lätt observerbarhet.

Jämförelsetabell

Tillvägagångssätt Bounded Spårning Per nyckel Signaler Självrengörande komplexitet |----------|:-------:|:--------:|:-------:|:-------:|:-------------:|:----------:| | Parallel.ForEachAsync Obligatoriska uppgifter som krävs för att uppfylla kraven i bilaga II till förordning (EU) nr 1094/2010 Dataflöde för dataflöde för TPL till Kanalerna och . . . . . . . och . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . till Polly på grund av att den inte uppfyller kraven i bilaga II till förordning (EU) nr 1094/2010. på bakgrundstjänster . . . . . . . . . och . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . MångfaldTransit/NServiceBuss på ett eller annat sätt | Efemeralt bibliotek Om du inte är säker på att det inte finns någon anledning att tro att det inte finns någon anledning att tro att du inte är det.


EfemeralForEachAsync: Den en-heta versionen

Från ParallellEfemeral.cs:

// Simple parallel processing with tracking
await items.EphemeralForEachAsync(
    async (item, ct) => await ProcessAsync(item, ct),
    new EphemeralOptions { MaxConcurrency = 8 });

// With keyed execution (per-user sequential)
await commands.EphemeralForEachAsync(
    cmd => cmd.UserId,  // Key selector
    async (cmd, ct) => await ExecuteCommandAsync(cmd, ct),
    new EphemeralOptions
    {
        MaxConcurrency = 32,
        MaxConcurrencyPerKey = 1  // Sequential per user
    });

Varför keyed pipelines materia

Tänk dig att bearbeta användarkommandon:

  • Användare A skickar kommandon 1, 2, 3
  • Användare B skickar kommandon 4, 5, 6

Utan nyckeltagning kan dessa verkställas enligt följande: 1, 4, 2, 5, 3, 6 - interleaved.

med MaxConcurrencyPerKey = 1:

  • Användaren A: s kommandon kör i ordning: 1 → 2 → 3
  • Användaren B kommandon köra i ordning: 4 → 5 → 6
  • Men A och B kan köra parallellt

Det här är per-enhet sekventiell, globalt parallell - kritiska för system där ordern är viktig inom en enhet.


Arbetskoordinatorn: En långlivad kö

Från Efemärarbetssamordnare.cs:

await using var coordinator = new EphemeralWorkCoordinator<TranslationRequest>(
    async (request, ct) => await TranslateAsync(request, ct),
    new EphemeralOptions
    {
        MaxConcurrency = 8,
        MaxTrackedOperations = 500,
        EnableDynamicConcurrency = true  // Allow runtime adjustment
    });

// Enqueue items over time
await coordinator.EnqueueAsync(new TranslationRequest("Hello", "es"));

// Check status anytime
Console.WriteLine($"Pending: {coordinator.PendingCount}");
Console.WriteLine($"Active: {coordinator.ActiveCount}");

// Get snapshots
var snapshot = coordinator.GetSnapshot();
var running = coordinator.GetRunning();
var failed = coordinator.GetFailed();
var completed = coordinator.GetCompleted();

// Control flow
coordinator.Pause();   // Stop pulling new work
coordinator.Resume();  // Continue

// Adjust concurrency at runtime (requires EnableDynamicConcurrency)
coordinator.SetMaxConcurrency(16);

// Pin important operations to survive eviction
coordinator.Pin(operationId);
coordinator.Unpin(operationId);
coordinator.Evict(operationId);

// When done
coordinator.Complete();
await coordinator.DrainAsync();

Kontinuerliga strömmar med IAsyncEnumerabel

await using var coordinator = EphemeralWorkCoordinator<Message>.FromAsyncEnumerable(
    messageStream,  // IAsyncEnumerable<Message>
    async (msg, ct) => await ProcessMessageAsync(msg, ct),
    new EphemeralOptions { MaxConcurrency = 16 });

await coordinator.DrainAsync();

Nyckelsamordnaren: Per-Entity Pipelines

Från EfemeralKeyedarbetssamordnare.cs:

await using var coordinator = new EphemeralKeyedWorkCoordinator<string, Command>(
    cmd => cmd.UserId,  // Key selector
    async (cmd, ct) => await ExecuteCommandAsync(cmd, ct),
    new EphemeralOptions
    {
        MaxConcurrency = 32,
        MaxConcurrencyPerKey = 1,      // Per-user sequential
        EnableFairScheduling = true,   // Prevent hot user starvation
        FairSchedulingThreshold = 10   // Reject if user has 10+ pending
    });

// TryEnqueue returns false if fair scheduling rejects
if (!coordinator.TryEnqueue(hotUserCommand))
{
    await DeferCommandAsync(hotUserCommand);
}

// Per-key visibility
var pendingForUser = coordinator.GetPendingCountForKey("user-123");
var opsForUser = coordinator.GetSnapshotForKey("user-123");

Resultatsamordnare

Från EfemeralResultatsamordnare.cs:

await using var coordinator = new EphemeralResultCoordinator<SessionInput, SessionResult>(
    async (input, ct) =>
    {
        var fingerprint = await ComputeFingerprintAsync(input.Events, ct);
        return new SessionResult(fingerprint, input.Events.Length);
    },
    new EphemeralOptions { MaxConcurrency = 16 });

await coordinator.EnqueueAsync(session);
coordinator.Complete();
await coordinator.DrainAsync();

// Get just the results (no metadata)
var results = coordinator.GetResults();

// Get snapshots with results + metadata
var snapshots = coordinator.GetSnapshot();

// Get base snapshots without results (privacy-safe)
var baseSnapshots = coordinator.GetBaseSnapshot();

// Filter by success/failure
var successful = coordinator.GetSuccessful();
var failed = coordinator.GetFailed();

Koncurrenskontroll

Från KoncurrencyGates.cs:

Biblioteket tillhandahåller två mekanismer för konvergenskontroll:

FastkoncurrencyGate (förvalt)

  • Stödd av SemaphoreSlim
  • Optimal värmeledningsprestanda
  • Kan inte justeras vid körning

Inställbar konvergensGate

  • Anpassat genomförande med Queue<WaiterEntry>
  • Stöd UpdateLimit() vid körning
  • Aktiverad via EnableDynamicConcurrency = true
// Dynamic concurrency adjustment
var coordinator = new EphemeralWorkCoordinator<T>(body,
    new EphemeralOptions
    {
        MaxConcurrency = 4,
        EnableDynamicConcurrency = true
    });

// Later, based on system load:
coordinator.SetMaxConcurrency(16);  // Scale up
coordinator.SetMaxConcurrency(2);   // Scale down

Fabrikens mönster: Namngivna samordnare

Från BeroendeInjicering.cs:

Så här är det. IHttpClientFactory, Du kan registrera namngivna konfigurationer:

// Registration
services.AddEphemeralWorkCoordinator<TranslationRequest>("fast",
    async (request, ct) => await FastTranslateAsync(request, ct),
    new EphemeralOptions { MaxConcurrency = 32 });

services.AddEphemeralWorkCoordinator<TranslationRequest>("accurate",
    async (request, ct) => await AccurateTranslateAsync(request, ct),
    new EphemeralOptions { MaxConcurrency = 4 });

// Usage
public class TranslationService(IEphemeralCoordinatorFactory<TranslationRequest> factory)
{
    private readonly EphemeralWorkCoordinator<TranslationRequest> _fast =
        factory.CreateCoordinator("fast");
    private readonly EphemeralWorkCoordinator<TranslationRequest> _accurate =
        factory.CreateCoordinator("accurate");
}

Fabriksgarantier

  1. Samma namn = samma instans - Jag ringer. CreateCoordinator("fast") två gånger returnerar samma samordnare
  2. Olika namn = olika instanser - "fast" och "accurate" få separata koordinatorer
  3. Löjligt skapande - Samordnare skapas först när de begärs
  4. Inställningsvalidering - Begära ett oregistrerat namn kastar ett användbart fel

Signalförfrågans API@ info: whatsthis

Alla samordnare tillhandahåller optimerade signalförfrågningsmetoder:

// Get all signals
var signals = coordinator.GetSignals();

// Filter by key (zero-allocation)
var userSignals = coordinator.GetSignalsByKey("user-123");

// Filter by time range
var recentSignals = coordinator.GetSignalsSince(DateTimeOffset.UtcNow.AddMinutes(-5));
var rangeSignals = coordinator.GetSignalsByTimeRange(from, to);

// Filter by signal name or pattern
var rateSignals = coordinator.GetSignalsByName("rate-limit");
var httpSignals = coordinator.GetSignalsByPattern("http.*");

// Check existence (short-circuits on first match)
if (coordinator.HasSignal("rate-limit"))
    await ThrottleAsync();

if (coordinator.HasSignalMatching("error.*"))
    await AlertAsync();

// Count signals efficiently (no allocation)
var totalSignals = coordinator.CountSignals();
var errorCount = coordinator.CountSignals("error");
var httpCount = coordinator.CountSignalsMatching("http.*");

Produktionsoptimeringar

Snabb ID-generering

Från EfemeralIdGenerator.cs:

internal static class EphemeralIdGenerator
{
    private static long _counter;
    private static readonly long _processStart = Environment.TickCount64;
    private static readonly int _processId = Environment.ProcessId;

    [MethodImpl(MethodImplOptions.AggressiveInlining)]
    public static long NextId()
    {
        var counter = Interlocked.Increment(ref _counter);

        // Combine counter with process-unique seed
        Span<byte> buffer = stackalloc byte[24];
        BitConverter.TryWriteBytes(buffer, _processStart);
        BitConverter.TryWriteBytes(buffer.Slice(8), _processId);
        BitConverter.TryWriteBytes(buffer.Slice(16), counter);

        return unchecked((long)XxHash64.HashToUInt64(buffer));
    }
}
  • Fritt från tilldelning (användningar) stackalloc)
  • Trådsäker (användningar) Interlocked.Increment)
  • Unikt för alla processer (inkluderar process-ID)
  • Icke-löpande (hash sprider disken)

Minnessäker långlivad drift

Koordinatorerna lagrar inte Task referenser - bara räknare:

private int _activeTaskCount;
private readonly TaskCompletionSource _drainTcs;

// In ExecuteItemAsync:
finally
{
    // Signal drain when last task completes AND channel iteration is done
    if (Interlocked.Decrement(ref _activeTaskCount) == 0 &&
        Volatile.Read(ref _channelIterationComplete))
    {
        _drainTcs.TrySetResult();
    }
}

Rengöring av per nyckel

Den nyckelfärdiga koordinatorn rensar automatiskt upp sysslolösa per nyckelsemaforer:

private sealed class KeyLock(SemaphoreSlim gate, int maxCount)
{
    public SemaphoreSlim Gate { get; } = gate;
    public int MaxCount { get; } = maxCount;
    public long LastUsedTicks = Environment.TickCount64;
}

// Cleanup runs periodically, removes locks idle > 60 seconds

Fullständigt exempel

// Program.cs
var builder = WebApplication.CreateBuilder(args);

// Named coordinators
builder.Services.AddEphemeralWorkCoordinator<TranslationRequest>("fast",
    async (req, ct) => await FastTranslateAsync(req, ct),
    new EphemeralOptions { MaxConcurrency = 16 });

// Keyed coordinator for per-user commands
builder.Services.AddEphemeralKeyedWorkCoordinator<string, UserCommand>("commands",
    cmd => cmd.UserId,
    sp =>
    {
        var handler = sp.GetRequiredService<ICommandHandler>();
        return async (cmd, ct) => await handler.HandleAsync(cmd, ct);
    },
    new EphemeralOptions
    {
        MaxConcurrency = 32,
        MaxConcurrencyPerKey = 1,
        EnableFairScheduling = true,
        CancelOnSignals = new HashSet<string> { "system-overload" }
    });

var app = builder.Build();
// Controller
[ApiController]
[Route("api")]
public class WorkController : ControllerBase
{
    private readonly EphemeralWorkCoordinator<TranslationRequest> _translator;
    private readonly EphemeralKeyedWorkCoordinator<string, UserCommand> _commands;

    public WorkController(
        IEphemeralCoordinatorFactory<TranslationRequest> translationFactory,
        IEphemeralKeyedCoordinatorFactory<string, UserCommand> commandFactory)
    {
        _translator = translationFactory.CreateCoordinator("fast");
        _commands = commandFactory.CreateCoordinator("commands");
    }

    [HttpPost("translate")]
    public async Task<IActionResult> Translate([FromBody] TranslationRequest request)
    {
        await _translator.EnqueueAsync(request);
        return Ok(new { pending = _translator.PendingCount });
    }

    [HttpPost("command")]
    public IActionResult SubmitCommand([FromBody] UserCommand command)
    {
        if (!_commands.TryEnqueue(command))
            return StatusCode(429, "Too many pending commands for this user");
        return Ok();
    }

    [HttpGet("status")]
    public IActionResult GetStatus() => Ok(new
    {
        translator = new
        {
            pending = _translator.PendingCount,
            active = _translator.ActiveCount,
            completed = _translator.TotalCompleted,
            failed = _translator.TotalFailed,
            hasRateLimit = _translator.HasSignal("rate-limit")
        },
        commands = new
        {
            pending = _commands.PendingCount,
            active = _commands.ActiveCount,
            errorCount = _commands.CountSignalsMatching("error.*")
        }
    });
}

Slutsatser

Vi har byggt ett komplett efemeralt avrättningsbibliotek med:

  1. EphemeralForEachAsync - En-shot parallell bearbetning med spårning
  2. EphemeralWorkCoordinator - Långlivade observerbara köer
  3. EphemeralKeyedWorkCoordinator - Per-entitet sekventiellt utförande med rättvis schemaläggning
  4. EphemeralResultCoordinator - Resultat-kaptur variant
  5. Fabrikens mönster - Namngivna konfigurationer som IHttpClientFactory
  6. Dynamisk överensstämmelse - Runtime justering av parallellism
  7. Signalinfrastruktur - Inbyggd signalemission och förfrågan

Mönstret sitter på en söt plats:

  • Mer observerbar än Parallel.ForEachAsync
  • Enklare än TPL Dataflöde
  • Mer integrerade än råa kanaler
  • Sekretessskydd genom konstruktion

Eld... och glöm inte.


Länkar

logo

© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.