Back to "Construcción de una biblioteca de ejecución efímera reutilizable"

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

Construcción de una biblioteca de ejecución efímera reutilizable

Friday, 12 December 2025

In Parte 1: Fuego y no lo hagas Bastante. Olvídelo., exploramos la teoría detrás de la ejecución efímera - limitado, privado, flujo de trabajo async desbuggable que recuerda lo suficiente para ser útil y luego evaporarse.

Este artículo convierte ese patrón en una biblioteca reutilizable que puede colocar en cualquier proyecto .NET.

Esto está ahora en el paquete Nuget mayoritariamente lucid.efemerals también más de 20 patrones y 'átomos' principalmente lucid.efemerals.

NuGet Licencia

Archivos de origen

La biblioteca se divide en archivos bien elaborados:

Archivo # # Propósito

|------|---------| | EfímeroOpciones.cs Configuración (condición, tamaño de la ventana, vida útil, señales) | EfímeroOperación.cs Seguimiento interno de la operación con soporte de señal | Snapshots.cs Registros de instantáneas inmutables expuestos a los consumidores | Señales.cs Los eventos de la señal, la propagación, las restricciones, y la señal global | EphemeralIdGenerator.cs Generación rápida de ID basada en XxHash64 | ConditionGates.cs Limitación de concurrencia fija y ajustable | StringPatternMatcher.cs Patrón estilo Glob emparejamiento para el filtrado de señales | ParaleloEfímero.cs Métodos de extensión estática (EphemeralForEachAsync) | | EphemeralWorkCoordinator.cs Coordinador de colas de trabajo de larga duración | EphemeralKeyedWorkCoordinator.cs Ejecución secuencial por llave con programación justa | EfímeroResultadoCoordinador.cs Variante del coordinador de captura de resultados | SeñalDispatcher.cs Enrutamiento de la señal de Async con la coincidencia del patrón | DependenciaInyección.cs métodos de extensión DI e implementaciones de fábrica | Ejemplos/SeñalizaciónHttpClient.cs Muestra de emisión de señal de grano fino para llamadas HTTP

Y Ensayos completos cubrir todas las cajas de bordes.


Antes y después

Esto es lo que estamos reemplazando:

// ❌ 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...

Y lo que estamos construyendo:

// ✅ 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();

La misma ejecución async. Observabilidad completa. No se conservan datos de usuario.


Inicio rápido

El patrón más común - registrar un coordinador en DI e inyectarlo:

// 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
    };
}

¿Qué variante necesito?

┌─────────────────────────────────────────────────────────────────┐
│                    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│
│                                                                 │
└─────────────────────────────────────────────────────────────────┘

El objeto de configuración

Desde EfímeroOpciones.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; }
}

Decisiones clave sobre el diseño

  • MaxConmoneda por defecto a la cuenta de CPU - sensible para el trabajo conectado a la CPU. Para el trabajo unido a E/S, aumentarlo.
  • HabilitarDinamicConmoneda permite ajustar el tiempo de funcionamiento a través de SetMaxConcurrency() - utiliza una puerta personalizada en lugar de SemaphoreSlim.
  • CancelarOnSeñals/DeferOnSeñals hacer que los coordinadores de la señal-reactiva - responden al estado del sistema ambiente (patrón de compatibilidad de soportes */?/listas de comas).
  • OnSignal es sincrónico; para el uso de fan-out de async SignalDispatcher o AsyncSignalProcessor dentro del controlador.
  • SignalConstraints evita bucles de señal infinitos con detección de ciclo y límites de profundidad.

Los registros de la instantánea

Desde 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);

Esto es sólo metadatos. Note lo que es no aquí:

  • Sin carga útil
  • Sin datos de entrada
  • Sin contenido de usuario

Lo suficiente para responder "¿qué pasó, cuándo, y funcionó?" - nada más.


Cómo se compara esto con otros enfoques

.NET le da varias maneras de hacer trabajo paralelo. He aquí cómo se compara la biblioteca efímera:

Paralelo.ParaEachAsync (.NET 6+)

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

Lo mejor para: Procesamiento simple paralelo de colecciones donde no se necesita visibilidad.

Lo que le falta:

  • Sin seguimiento de operaciones
  • Sin ejecución secuencial por tecla
  • No hay visibilidad en lo que está corriendo

Use Efemeral cuando: Usted necesita depuración/observabilidad, pedido por llave, o procesamiento de señal-reactiva.

Flujo de datos TPL

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;

Lo mejor para: Tuberías complejas de flujo de datos con ramificación, fusión, loteado.

Lo que hace bien:

  • Composición rica de tuberías (enlace bloques juntos)
  • Lotes incorporados, transformándose, transmitiendo
  • Capacidad limitada con contrapresión

Usar el flujo de datos TPL cuando: Se necesitan topologías de tuberías complejas (fan-out, fan-in, enrutamiento condicional).

Use Efemeral cuando: Necesita seguimiento de operaciones, API más simple, o coordinación de señal-reactiva.

Sistema.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);

Lo mejor para: Patrones productor-consumidor donde se controla a ambos lados.

Lo que hace bien:

  • Excelente rendimiento
  • Retraso de la presión a través de canales delimitados
  • Separación de productores y consumidores

Usar canales cuando: Está construyendo una infraestructura personalizada y necesita el máximo control.

Use Efemeral cuando: Usted quiere seguimiento de la operación y la observabilidad sin la placa de caldera.

Polly

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

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

Lo mejor para: Políticas de resiliencia (intensidad, disyuntor, tiempo de espera) para operaciones individuales.

Usar Polly cuando: Necesitas resiliencia en torno a llamadas individuales.

Use Efemeral cuando: Usted necesita coordinación a través de muchas operaciones con la conciencia ambiental.

Combínalos.: Utilice Polly dentro de su cuerpo de trabajo efímero para la resiliencia per-operativa.

MassTransit / NServiceBus

Lo mejor para: Mensajería distribuida entre servicios con colas duraderas.

Usar buses de mensajes cuando: El trabajo debe sobrevivir a los reinicios del proceso, abarcar múltiples servicios o requerir entrega garantizada.

Use Efemeral cuando: El trabajo está en proceso, no necesita durabilidad, y usted quiere la observabilidad ligera.

Cuadro comparativo

Aproximación Atado Seguimiento Per-Key Señales Auto-Limpieza Complejidad |----------|:-------:|:--------:|:-------:|:-------:|:-------------:|:----------:| | Parallel.ForEachAsync # # # # # # # # # # N/A # # Baja # # TPL Flujo de datos â â € â € â € â € â € â € TM Alto â € â € TM Alto

Canales # # Canales # # Canales # # Canal # # Canal # # Canal # # Canal # # Canal # # Canal # # Canal # # Canal # # Canal # # Canal # # Canal # Canal # # Canal # # Canal # # Canal # # Canal # # Canal # # Canal # # Canal # # Canal # # Canal # # Canal # # Canal # # Canal # # Canal # # Canal # # Canal # # Canal # # # Canal # # Canal # # # Canal # # # # Canal # # # Canal # # # # Canal # # Canal # # Canal # # Canal # # Media # # Media # Media # # Media

Polly # # N/A # # N/A # # N/A # Baja

Servicios de antecedentes # # # # # # # Mediano

MassTransit/NServiceBus # # # # # MassTransit/NServiceBus # # # High #

| Biblioteca efímera # Low # # Low # # Low # # Low # # Low # # Low # # Low # # Low # # Low # # Low # # Low # # Low # # Low # # Low # # Low # # Low # # Low # # Low # # Low # # Low # # Low # # Low #


EfímeroParaCadaAsync: La versión única

Desde ParaleloEfímero.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
    });

Por qué importan los oleoductos con llave

Imagine procesar comandos de usuario:

  • El usuario A envía comandos 1, 2, 3
  • El usuario B envía comandos 4, 5, 6

Sin keying, éstos podrían ejecutarse como: 1, 4, 2, 5, 3, 6 - entrelazados.

Con MaxConcurrencyPerKey = 1:

  • Los comandos del usuario A se ejecutan en orden: 1 → 2 → 3
  • Los comandos del usuario B se ejecutan en orden: 4 → 5 → 6
  • Pero A y B pueden correr en paralelo

Esto es secuencial por entidad, paralelo a nivel mundial - críticos para los sistemas en los que el orden importa dentro de una entidad.


El coordinador del trabajo: una cola de larga duración

Desde EphemeralWorkCoordinator.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();

Corrientes continuas con IAsyncEnumerable

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

await coordinator.DrainAsync();

El coordinador clave: los oleoductos Per-Entity

Desde EphemeralKeyedWorkCoordinator.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");

Coordinadores de captura de resultados

Desde EfímeroResultadoCoordinador.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();

Control de las monedas

Desde ConditionGates.cs:

La biblioteca ofrece dos mecanismos de control de las monedas:

ConditionGate Fijo (Default)

  • Respaldado por SemaphoreSlim
  • Rendimiento óptimo del hot-path
  • No se puede ajustar en tiempo de ejecución

Puerta de Condición Ajustable

  • Aplicación personalizada con Queue<WaiterEntry>
  • Apoyos UpdateLimit() en tiempo de ejecución
  • Activado a través de 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

El Patrón de Fábrica: Nombrados Coordinadores

Desde DependenciaInyección.cs:

Me gusta IHttpClientFactory, puede registrar configuraciones con nombre:

// 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");
}

Garantías de fábrica

  1. Mismo nombre = misma instancia - Llamando. CreateCoordinator("fast") dos veces devuelve el mismo coordinador
  2. Distintos nombres = diferentes instancias - "fast" y "accurate" conseguir coordinadores separados
  3. Creación perezosa - Los coordinadores sólo se crean cuando se solicitan por primera vez
  4. Validación de la configuración - Solicitar un nombre no registrado arroja un error útil

API de búsqueda de señales

Todos los coordinadores proporcionan métodos optimizados de consulta de señales:

// 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.*");

Optimización de la producción

Generación de ID rápida

Desde EphemeralIdGenerator.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));
    }
}
  • Libre de asignaciones (usos stackalloc)
  • A prueba de hilos (usos Interlocked.Increment)
  • Único en todos los procesos (incluye ID del proceso)
  • No secuenciales (hash difunde el contador)

Operación de larga vida segura para la memoria

Los coordinadores no almacenan Task referencias - sólo contadores:

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();
    }
}

Limpieza de bloqueo por clave

El coordinador con llave limpia automáticamente los semáforos inactivos por tecla:

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

Ejemplo completo

// 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.*")
        }
    });
}

Conclusión

Hemos construido una completa biblioteca de ejecución efímera con:

  1. EphemeralForEachAsync - Procesamiento paralelo de un solo disparo con seguimiento
  2. EphemeralWorkCoordinator - Colas observables de larga duración
  3. EphemeralKeyedWorkCoordinator - Ejecución secuencial por entidad con programación justa
  4. EphemeralResultCoordinator - Variante de captura de resultados
  5. Patrón de fábrica - Configuraciones nombradas como IHttpClientFactory
  6. Condición dinámica - Ajuste de tiempo de ejecución del paralelismo
  7. Infraestructura de señalización - Emisión de señal incorporada y consulta

El patrón se sienta en un punto dulce:

  • Más observable que Parallel.ForEachAsync
  • Más simple que TPL Dataflow
  • Más integrados que los canales crudos
  • Privacidad segura por diseño

Fuego... y no lo olvides.


Vínculos

logo

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