mayormente lucid.efemeral.complete; un extraño sistema de patrón de sistemas concurrentes en una caché LRU. (Español (Spanish))

mayormente lucid.efemeral.complete; un extraño sistema de patrón de sistemas concurrentes en una caché LRU.

Sunday, 14 December 2025

//

24 minute read

Bueno, esta ha sido mi obsesión durante la semana pasada. Vea las partes anteriores y lo que llevó a esto; '¿Qué pasa si un LRU era un contexto de ejecución.'. Ahora es un conjunto de 30 paquetes Nuget que cubren la mayoría de los principales patrones de ejecución concurrentes (en paquetes de línea TINY 5-10). Obtenga capacidades adaptativas increíbles con una sintaxis SIMPLE!

Puede encontrar la fuente aquí: https://github.com/scottgal/mostlylucid.atoms/blob/main/mostlylucid.efemeral/src/mostlylucid.efemeral.complete

Prioridades

Lea la sección [Parte anterior sobre Señales ]Presentado aquí es el Readme.md de la pacakage mayoritariamentelucid.efemeral.complete que contiene tanto el núcleo principalmentelucid.efemeral paquete y todos los patrones, 'átomos' (coordinadores, etc) paquetes en un DLL conveniente.

Recursos básicos

O utilice el núcleo mayormente lúcido.efemeral un TINY (literalmente 10 clases) que le da toda la funcionalidad en bruto.

Atributos & DI

O si desea un enrutamiento async basado en atributo completo con sencillo [EphemeralJob] y el servicio.El estilo de registro del coordinador utiliza el sobre todo lucid.efemeral.attributes paquete.

Esto es probablemente el tema de mi blog que va hacia adelante ... usted ha sido advertido

Mayormente lúcido.Efímero.Completa

NuGet

Todo de Mostlylucid.Ephemeral en una sola DLL - ejecución async limitada con coordinación basada en la señal.

dotnet add package mostlylucid.ephemeral.complete

Este paquete compila todo el núcleo, átomo y código de patrón en un solo conjunto. Para paquetes individuales, vea los enlaces en cada uno sección que figura a continuación.


Índice


Inicio rápido

using Mostlylucid.Ephemeral;

// Long-lived work coordinator
await using var coordinator = new EphemeralWorkCoordinator<WorkItem>(
    async (item, ct) => await ProcessAsync(item, ct),
    new EphemeralOptions { MaxConcurrency = 8 });

await coordinator.EnqueueAsync(new WorkItem("data"));

// One-shot parallel processing
await items.EphemeralForEachAsync(
    async (item, ct) => await ProcessAsync(item, ct),
    new EphemeralOptions { MaxConcurrency = 8 });

Inscripción de los servicios

var builder = WebApplication.CreateBuilder(args);

builder.Services.AddCoordinator<WorkItem>(
    async (item, ct) => await ProcessAsync(item, ct),
    new EphemeralOptions { MaxConcurrency = 8, MaxTrackedOperations = 128 });

builder.Services.AddEphemeralSignalJobRunner<LogWatcherJobs>();

var app = builder.Build();
app.MapPost("/", async ([FromServices] IEphemeralCoordinatorFactory<WorkItem> factory, WorkItem item) =>
{
    var coordinator = factory.CreateCoordinator();
    await coordinator.EnqueueAsync(item);
    return Results.Accepted();
});

await app.RunAsync();

Los familiares services.AddCoordinator<T>() Ayudantes y AddEphemeralSignalJobRunner<T>() mantenga el registro de servicio conciso, deje que DI sea el dueño del fregadero/corredor y haga que las nuevas historias de responsabilidad/caché/loging estén a un solo clic.

Empleos basados en atributos

mostlylucid.ephemeral.complete bultos mostlylucid.ephemeral.attributes, así que las tuberías de atributo son parte del núcleo superficie. Trate al corredor como un consumidor de señal de primera clase: métodos decorados se unen a la misma caché, registro, y historias de fijación, y cada atributo puede declarar Priority, nivel de empleo MaxConcurrency, Lane, Key fuentes, señal las emisiones, las anulaciones de pin/expire y los reintentos.

Pomos de atributo clave:

  • Ordenar & carriles: Usar Priority, MaxConcurrency, y Lane para mantener el trabajo en orden determinista mientras los caminos calientes Mantente separado.
  • Teclado & etiquetado: OperationKey, KeyFromSignal, KeyFromPayload, y [KeySource] ayudarle a trabajar en grupo con claves significativas para el registro, la programación justa, y diagnósticos.
  • Fijación y reintentos: Pin, ExpireAfterMs, AwaitSignals, MaxRetries, y RetryDelayMs dejar que los manipuladores se extiendan su visibilidad, ejecución de la puerta hasta que las dependencias lleguen, y sanar con reintentos mientras emiten señales de fallo.
  • Coreografía de señales: Emit EmitOnStart, EmitOnComplete, y EmitOnFailure para señalizar las etapas posteriores, vigilantes, u otros coordinadores sin cableado manual.
var sink = new SignalSink();
await using var runner = new EphemeralSignalJobRunner(sink, new[] { new LogWatcherJobs(sink) });

var loggerFactory = LoggerFactory.Create(builder =>
{
    builder.AddConsole();
    builder.AddProvider(new SignalLoggerProvider(new TypedSignalSink<SignalLogPayload>(sink)));
});

var logger = loggerFactory.CreateLogger("orders");
logger.LogError(new EventId(1001, "DbFailure"), "Order store failed");

// Later tasks or other services can also raise watcher-friendly signals directly:
sink.Raise("log.error.orders.dbfailure", key: "orders");
public sealed class LogWatcherJobs
{
    private readonly SignalSink _sink;

    public LogWatcherJobs(SignalSink sink) => _sink = sink;

    [EphemeralJob("log.error.*", Priority = 1, MaxConcurrency = 2, Lane = "hot:4", EmitOnComplete = new[] { "incident.created" })]
    public Task EscalateAsync(SignalEvent signal)
    {
        Console.WriteLine($"escalating {signal.Signal} for {signal.Key}");
        _sink.Raise("incident.created", key: signal.Key);
        return Task.CompletedTask;
    }

    [EphemeralJob("incident.created", EmitOnStart = new[] { "incident.monitor.start" })]
    public Task NotifyAsync(SignalEvent signal)
    {
        Console.WriteLine($"notified incident for {signal.Key}");
        return Task.CompletedTask;
    }
}

Este corredor ahora se sienta al inicio y reacciona cuando quiera log.error.* o cualquier señal emitida golpea el fregadero. Atributo los manipuladores también pueden leer las claves de las señales/cargas de carga, trabajar el pin hasta acks aguas abajo, emitir señales de terminación/fracaso, y ranura en carriles para ordenar. Para el uso de configuraciones DI-first services.AddEphemeralSignalJobRunner<T>() (o el ámbito de aplicación Variante) por lo que el corredor y el fregadero son gestionados por el contenedor.

[EfemeralJobs(SignalPrefix = "etapa", DefaultLane = "pipeline")] clase pública sellada StageJobs { [EfímeroJob("ingest", EmitOnComplete = nuevo[{ "stage.ingest.done" }) Public Task IngestAsync(SignalEvent evt) => Consola.Out.WriteLineAsync(evt.Signal);

[EphemeralJob("finalize")]
public Task FinalizeAsync(SignalEvent evt) => Console.Out.WriteLineAsync("final stage");

}

var stageSink = nueva señalSink(); esperar usando var stageRunner = nuevo EphemeralSignalJobRunner(etapaSink, nuevo[] { nuevo StageJobs() }; stageSink.Raise("stage.ingest");

Los puestos de trabajo pesados pueden confiar en ResponsibilitySignalManager.PinUntilQueried (patrón de ack predeterminado) responsibility.ack.*) a mantener sus operaciones visibles hasta que un lector de aguas abajo obtenga la carga útil, mientras que OperationEchoMaker/ OperationEchoAtom persistir la corriente final de la señal para que los auditores o moléculas puedan todavía “gustar” el último estado incluso después de el átomo muere.

Tareas programadas

mostlylucid.ephemeral.complete también contiene mostlylucid.ephemeral.atoms.scheduledtasksDefinir cron o JSON horarios a través de ScheduledTaskDefinition (cron, señal, opcional key, payload, description, timeZone, format, runOnStartup, etc.), y dejar ScheduledTasksAtom encolar trabajo duradero a través de DurableTaskAtom. Cada trabajo programado eleva la señal configurada dentro de una ventana de coordinador, por lo que hereda la semántica de fijación, registro y responsabilidad mientras sus moléculas o tuberías de atributos responden a la onda de señal emitida.

Cada DurableTask lleva el horario Name, Signal, opcional Key, incluso un mecanografiado Payload, y Description, así que los oyentes posteriores saben inmediatamente qué trabajo se ha ejecutado y qué metadatos (nombres de archivo, URLs, etc.) consumir. DurableTaskAtom.WaitForIdleAsync() cuando usted sólo quiere esperar a que el estallido actual de trabajo programado para terminar sin completar el átomo, mantener el programador listo para la siguiente cron garrapata.

& Señales de registro

mostlylucid.ephemeral.logging espejos Microsoft.Extensiones.Entrando en señales y viceversa. SignalLoggerProvider a su fábrica de registradores para que los eventos de registro suban log.* señales, y gancho SignalToLoggerAdapter si desea que las señales fluyan de nuevo en la tubería de registro estándar.

var sink = new SignalSink();
var typedSink = new TypedSignalSink<SignalLogPayload>(sink);

using var loggerFactory = LoggerFactory.Create(builder =>
{
    builder.AddConsole();
    builder.AddProvider(new SignalLoggerProvider(typedSink));
});

using var watcher = new EphemeralSignalJobRunner(sink, new[] { new LogWatcherJobs(sink) });

var logger = loggerFactory.CreateLogger("orders");
logger.LogError(new EventId(1001, "DbFailure"), "Order store failed");
public sealed class LogWatcherJobs
{
    private readonly SignalSink _sink;

    public LogWatcherJobs(SignalSink sink) => _sink = sink;

    [EphemeralJob("log.error.*")]
    public Task EscalateAsync(SignalEvent signal)
    {
        _sink.Raise("incident.created", key: signal.Key);
        return Task.CompletedTask;
    }

    [EphemeralJob("incident.created")]
    public Task NotifyAsync(SignalEvent signal)
    {
        Console.WriteLine($"Incident for {signal.Key}");
        return Task.CompletedTask;
    }
}

Uso SignalToLoggerAdapter para reflejar las señales resultantes de nuevo en registros estándar para que la pila de monitoreo vea ambos Lados del puente.


Coordinadores básicos

Paquete: mayormente lúcido.efemeral

Coordinador de Trabajo Efímero<T>

Cola de trabajo de larga duración con concurrencia limitada y ventana observable.

await using var coordinator = new EphemeralWorkCoordinator<Request>(
    async (req, ct) => await HandleAsync(req, ct),
    new EphemeralOptions
    {
        MaxConcurrency = 8,
        MaxTrackedOperations = 200,
        MaxOperationLifetime = TimeSpan.FromMinutes(5)
    });

await coordinator.EnqueueAsync(request);

// Observe state
var running = coordinator.GetRunning();
var failed = coordinator.GetFailed();
var pending = coordinator.PendingCount;

// Graceful shutdown
coordinator.Complete();
await coordinator.DrainAsync();

EphemeralKeyedWorkCoordinator<TKey, T>

Procesamiento secuencial por llave - elementos con la misma clave procesada en orden.

await using var coordinator = new EphemeralKeyedWorkCoordinator<Order, string>(
    order => order.CustomerId,  // Key selector
    async (order, ct) => await ProcessOrder(order, ct),
    new EphemeralOptions
    {
        MaxConcurrency = 16,      // Total parallel
        MaxConcurrencyPerKey = 1  // Sequential per customer
    });

await coordinator.EnqueueAsync(order);

EfímeroResultadoCoordinador<TInput, TResult>

Captura los resultados de las operaciones asíncronas.

await using var coordinator = new EphemeralResultCoordinator<Request, Response>(
    async (req, ct) => await FetchAsync(req, ct),
    new EphemeralOptions { MaxConcurrency = 4 });

var id = await coordinator.EnqueueAsync(request);
var snapshot = await coordinator.WaitForResult(id);
if (snapshot.HasResult)
    Console.WriteLine(snapshot.Result);

Coordinador de trabajo prioritario<T>

Múltiples carriles prioritarios con concurrencia configurable por carril.

var coordinator = new PriorityWorkCoordinator<WorkItem>(
    async (item, ct) => await ProcessAsync(item, ct),
    new PriorityWorkCoordinatorOptions<WorkItem>(
        Lanes: new[] { new PriorityLane("high"), new PriorityLane("normal"), new PriorityLane("low") }
    ));

await coordinator.EnqueueAsync(item, "high");

Configuración (Opciones Efímeras)

new EphemeralOptions
{
    // Concurrency
    MaxConcurrency = 8,                    // Max parallel operations
    MaxConcurrencyPerKey = 1,              // For keyed coordinators
    EnableDynamicConcurrency = false,      // Allow runtime adjustment

    // Memory
    MaxTrackedOperations = 200,            // Window size (LRU eviction)
    MaxOperationLifetime = TimeSpan.FromMinutes(5),

    // Fair scheduling (keyed only)
    EnableFairScheduling = false,          // Prevent hot key starvation
    FairSchedulingThreshold = 10,

    // Signals
    Signals = sharedSink,                  // Shared signal sink
    OnSignal = evt => { },                 // Sync callback
    OnSignalAsync = async (evt, ct) => { }, // Async callback
    CancelOnSignals = new HashSet<string> { "circuit-open" },
    DeferOnSignals = new HashSet<string> { "backpressure" },
    DeferCheckInterval = TimeSpan.FromMilliseconds(100),
    MaxDeferAttempts = 50,

    // Signal handler limits
    MaxConcurrentSignalHandlers = 4,
    MaxQueuedSignals = 1000
}

Señales

Las operaciones emiten señales de observabilidad transversal.

// Query signals
bool hasError = coordinator.HasSignal("error");
int count = coordinator.CountSignals("error");
var errors = coordinator.GetSignalsByPattern("error.*");

// Shared sink across coordinators
var sink = new SignalSink();
var c1 = new EphemeralWorkCoordinator<A>(body, new EphemeralOptions { Signals = sink });
var c2 = new EphemeralWorkCoordinator<B>(body, new EphemeralOptions { Signals = sink });
sink.Raise("system.busy");  // Both see it

Señales de Responsabilidad y Finalización

¿Necesita mantener los resultados visibles el tiempo suficiente para los consumidores de aguas abajo? ResponsibilitySignalManager le permite fijar un operación hasta que llegue una señal de ack (patrón por defecto) responsibility.ack.* con key=operationId). Proporcionar un opcional description por lo que la operación puede describir su responsabilidad, y establecer maxPinDuration a la gracia auto-claro si el consumidor nunca aparece.

var manager = new ResponsibilitySignalManager(coordinator, sink, maxPinDuration: TimeSpan.FromMinutes(5));
if (manager.PinUntilQueried(operationId, "file.ready", ackKey: fileId, description: "Awaiting fetch"))
{
    sink.Raise("file.ready", key: fileId);
}
// Consumer acknowledges the work
sink.Raise("file.ready.ack", key: fileId);
using Mostlylucid.Ephemeral.Patterns;

var notes = new LastWordsNoteAtom(async note => await noteRepository.SaveAsync(note));
coordinator.OperationFinalized += snapshot =>
{
    var note = new LastWordsNote(
        OperationId: snapshot.OperationId,
        Key: snapshot.Key,
        Signal: snapshot.Signals?.FirstOrDefault(),
        Timestamp: DateTimeOffset.UtcNow);

    _ = notes.EnqueueAsync(note);
};

LastWordsNote permanece pequeño (id de operación, clave, señal, marca de tiempo), para que pueda registrar cualquier estado mínimo que le importe cerca de antes de que se recoja la operación.

El coordinador también mantiene un eco de corta duración de las señales finales (habilitado a través de EnableOperationEcho) que usted puede inspeccionar con GetEchoes() cuando necesite reproducir la onda de señal recortada sin mantener la operación completa alrededor.

var recentErrors = coordinator.GetEchoes(pattern: "error.*")
    .Where(e => e.Timestamp > DateTimeOffset.UtcNow - TimeSpan.FromMinutes(1))
    .ToList();

if (recentErrors.Any())
    logger.LogWarning("Trimmed errors: {Count}", recentErrors.Count);

OperationEchoRetention y OperationEchoCapacity Deja que equilibres cuántos ecos guardas y cuánto tiempo permanecen, para que pueda reproducir las “últimas palabras” el tiempo suficiente para los diagnósticos de superficie.

El administrador despunta automáticamente cuando el ack se dispara, pero puedes llamar CompleteResponsibility(operationId) para poner fin a la la responsabilidad temprana (por ejemplo, en los reintentos). OperationFinalized cuando la ventana los recorta, así que suscribirse si desea emitir una señal final, diagnósticos de registro, o ejecutar “últimas palabras” limpieza.


Átomos (bloques de construcción)

FixedWorkAtom

**Paquete: ** principalmente lúcido.efemeral.átomos.trabajo fijo

Grupo de trabajadores fijos con estadísticas. Envoltorio de API mínima alrededor de EphemeralWorkCoordinator.

using Mostlylucid.Ephemeral.Atoms.FixedWork;

await using var atom = new FixedWorkAtom<WorkItem>(
    async (item, ct) => await ProcessAsync(item, ct),
    maxConcurrency: 4,
    maxTracked: 200);

await atom.EnqueueAsync(item);

// Get stats
var (pending, active, completed, failed) = atom.Stats();
Console.WriteLine($"Completed: {completed}, Failed: {failed}");

// Get recent operations
var snapshot = atom.Snapshot();

// Graceful shutdown
await atom.DrainAsync();

KeyedSequentialAtom

**Paquete: ** mayormente lúcido.efemeral.átoms.keyedsequential

Procesamiento secuencial por llave con programación justa opcional.

using Mostlylucid.Ephemeral.Atoms.KeyedSequential;

await using var atom = new KeyedSequentialAtom<Order, string>(
    keySelector: order => order.CustomerId,
    body: async (order, ct) => await ProcessOrder(order, ct),
    maxConcurrency: 16,
    perKeyConcurrency: 1,           // Sequential per key
    enableFairScheduling: true);    // Prevent hot key starvation

await atom.EnqueueAsync(order1);  // Customer A
await atom.EnqueueAsync(order2);  // Customer A - waits for order1
await atom.EnqueueAsync(order3);  // Customer B - parallel with A

var (pending, active, completed, failed) = atom.Stats();
await atom.DrainAsync();

SeñalAwareAtom

**Paquete: ** principalmente lucid.efemeral.átoms.signalaware

Pausa o cancela la entrada en función de las señales ambientales.

using Mostlylucid.Ephemeral.Atoms.SignalAware;

var sink = new SignalSink();

await using var atom = new SignalAwareAtom<WorkItem>(
    async (item, ct) => await ProcessAsync(item, ct),
    cancelOn: new HashSet<string> { "shutdown", "circuit-open" },
    deferOn: new HashSet<string> { "backpressure.*" },
    deferInterval: TimeSpan.FromMilliseconds(100),
    maxDeferAttempts: 50,
    signals: sink,
    maxConcurrency: 8);

// Enqueue work
await atom.EnqueueAsync(item);

// Raise ambient signals
atom.Raise("backpressure.downstream");  // New items defer
sink.Raise("shutdown");                  // New items rejected (returns -1)

await atom.DrainAsync();

BatchingAtom

**Paquete: ** principalmente lucid.efemeral.átoms.batching

Recoger los artículos en lotes por tamaño o intervalo de tiempo.

using Mostlylucid.Ephemeral.Atoms.Batching;

await using var atom = new BatchingAtom<LogEntry>(
    onBatch: async (batch, ct) =>
    {
        Console.WriteLine($"Flushing {batch.Count} entries");
        await FlushToDatabase(batch, ct);
    },
    maxBatchSize: 100,
    flushInterval: TimeSpan.FromSeconds(5));

// Items are batched automatically
atom.Enqueue(new LogEntry("User logged in"));
atom.Enqueue(new LogEntry("Request received"));
// ... batch flushes when full OR after 5 seconds

RetryAtom

Paquete: mayormente lucid.efemeral.átoms.retry

Envoltorio de prueba de retroceso exponencial.

using Mostlylucid.Ephemeral.Atoms.Retry;

await using var atom = new RetryAtom<ApiRequest>(
    async (req, ct) => await CallExternalApi(req, ct),
    maxAttempts: 3,
    backoff: attempt => TimeSpan.FromMilliseconds(100 * Math.Pow(2, attempt)),
    maxConcurrency: 4);

// Automatically retries on failure with exponential backoff
// Attempt 1: immediate
// Attempt 2: 200ms delay
// Attempt 3: 400ms delay
await atom.EnqueueAsync(new ApiRequest("https://api.example.com"));

await atom.DrainAsync();

Átomos de almacenamiento de datos

Paquete: mayormente lúcido.efemeral.átomos.datos

Configuración compartida para átomos de almacenamiento (DataStorageConfig, IDataStorageAtom<TKey, TValue>) además de las convenciones de señales que impulsan los adaptadores de archivos, SQLite y PostgreSQL.

using Mostlylucid.Ephemeral.Atoms.Data;
using Mostlylucid.Ephemeral.Atoms.Data.File;

var sink = new SignalSink();
var config = new DataStorageConfig
{
    DatabaseName = "orders",
    SignalPrefix = "save.data",
    LoadSignalPrefix = "load.data",
    DeleteSignalPrefix = "delete.data",
    MaxConcurrency = 1
};

await using var storage = new FileDataStorageAtom<string, Order>(sink, config, "./orders");

storage.EnqueueSave("order-123", new Order { Id = "order-123", Total = 42.00m });
var loaded = await storage.LoadAsync("order-123");

Usar el mismo DataStorageConfig con Mostlylucid.Ephemeral.Atoms.Data.Sqlite o Mostlylucid.Ephemeral.Atoms.Data.Postgres implementaciones para la persistencia duradera, impulsada por la señal alimentada por SQLite/Postgres. saved.data.{dbname} señales para iniciar el trabajo aguas abajo mientras load.data.{dbname} dispara cachés hidratados.


MoleculeRunner & AtomTrigger

**Paquete: ** principalmente lúcido.efemeral.átomos.moléculas

Planos compuestos con MoleculeBlueprintBuilder le permiten definir los átomos (pago, inventario, envío, la notificación) que debe ejecutarse cuando una señal como order.placed Llega. MoleculeRunner escucha el gatillo , crea un patrón compartido MoleculeContext, y ejecuta cada paso mientras se suscribe a los eventos de inicio/completado. AtomTrigger cuando la señal de un átomo debe iniciar otro coordinador o molécula.

var sink = new SignalSink();
var blueprint = new MoleculeBlueprintBuilder("order", "order.placed")
    .AddAtom(async (ctx, ct) => await paymentCoordinator.EnqueueAsync(ctx.TriggerSignal.Key!, ct))
    .AddAtom(async (ctx, ct) =>
    {
        ctx.Raise("order.payment.complete", ctx.TriggerSignal.Key);
        await inventoryCoordinator.EnqueueAsync(ctx.TriggerSignal.Key!, ct);
    })
    .Build();

await using var runner = new MoleculeRunner(sink, new[] { blueprint }, serviceProvider);
using var trigger = new AtomTrigger(sink, "order.payment.complete", async (signal, ct) =>
{
    await notificationCoordinator.EnqueueAsync(signal.Key!, ct);
});

sink.Raise("order.placed", key: "order-42");

Los pasos de moléculas pueden generar señales adicionales (ctx.Raise("order.shipping.start")) por lo que el resto del sistema recoge el batuta.


DeslizamientoCacheAtom

**Paquete: ** principalmente lucid.efemeral.átoms.slidingcache

Caché con vencimiento deslizante - el acceso a un resultado restablece su TTL.

using Mostlylucid.Ephemeral.Atoms.SlidingCache;

await using var cache = new SlidingCacheAtom<string, UserProfile>(
    async (userId, ct) => await LoadUserProfileAsync(userId, ct),
    slidingExpiration: TimeSpan.FromMinutes(5),
    absoluteExpiration: TimeSpan.FromHours(1),
    maxSize: 1000);

// First call: computes and caches
var profile = await cache.GetOrComputeAsync("user-123");

// Second call within 5 minutes: returns cached, resets TTL
var cached = await cache.GetOrComputeAsync("user-123");

// Try get without computation (still resets TTL on hit)
if (cache.TryGet("user-123", out var profile))
    Console.WriteLine(profile.Name);

// Get stats
var stats = cache.GetStats();
Console.WriteLine($"Entries: {stats.TotalEntries}, Hot: {stats.HotEntries}");

EfímeroLruCache

Paquete: núcleo (mostlylucid.ephemeral) — caché auto-optimizante con TTL deslizante en cada golpe y TTL extendido para Las llaves calientes.

using Mostlylucid.Ephemeral;

var cache = new EphemeralLruCache<string, Widget>(new EphemeralLruCacheOptions
{
    DefaultTtl = TimeSpan.FromMinutes(5),
    HotKeyExtension = TimeSpan.FromMinutes(30),
    HotAccessThreshold = 3,
    MaxSize = 10_000,
    SampleRate = 5 // emit 1 in 5 signals
});

var widget = await cache.GetOrAddAsync("widget:42", async key =>
{
    var data = await LoadWidgetAsync(key);
    return data!;
});

// Stats and signals to see how the cache self-focuses on hot keys
var stats = cache.GetStats();              // hot/expired counts, size
var signals = cache.GetSignals("cache.*"); // cache.hot/evict/miss/hit

Consejo: MemoryCache se puede configurar para la expiración deslizante, pero nunca emite las señales de calor/frío o extiende TTL para las llaves calientes. EphemeralLruCache es el valor predeterminado de auto-optimización en el paquete principal (y en SqliteSingleWriter) siempre que desee que la caché se centre en el conjunto de trabajo activo.

Eco Maker

Paquete: mayormente lúcido.efemeral.atoms.echo

Capture las “últimas palabras” mecanografiadas que emite una operación antes de recortarla. El átomo mantiene una ventana de señal limitada carga útil (combinación de ActivationSignalPattern / CaptureSignalPattern) y cuándo OperationFinalized fuegos que produce OperationEchoEntry<TPayload> registros que puede persistir a través de OperationEchoAtom<TPayload>.

var sink = new SignalSink();
var typedSink = new TypedSignalSink<EchoPayload>(sink);
var echoAtom = new OperationEchoAtom<EchoPayload>(async echo => await repository.AppendAsync(echo));

await using var coordinator = new EphemeralWorkCoordinator<JobItem>(ProcessAsync);
using var maker = coordinator.EnableOperationEchoing(
    typedSink,
    echoAtom,
    new OperationEchoMakerOptions<EchoPayload>
    {
        ActivationSignalPattern = "echo.capture",
        CaptureSignalPattern = "echo.*",
        MaxTrackedOperations = 128
    });

typedSink.Raise("echo.capture", new EchoPayload("order-1", "archived"), key: "order-1");

Los trabajos de atributos simplemente elevan la señal mecanografiada con cualquier estado que consideren crítico, y el fabricante mantiene el conjunto de trabajo limitada mientras persistes en el eco.


Patrones (Listos para usar)

Interruptor de circuitos basado en señales

**Paquete: ** principalmente lucid.efemeral.patterns.circuitbreaker

Disyuntor apátrida usando la ventana del historial de señales.

using Mostlylucid.Ephemeral.Patterns.CircuitBreaker;

var breaker = new SignalBasedCircuitBreaker(
    failureSignal: "api.failure",
    threshold: 5,
    windowSize: TimeSpan.FromSeconds(30));

// Check before making calls
if (breaker.IsOpen(coordinator))
{
    var retryAfter = breaker.GetTimeUntilClose(coordinator);
    throw new CircuitOpenException("Too many failures", retryAfter);
}

// Pattern matching variant
if (breaker.IsOpenMatching(coordinator, "error.*"))
    throw new CircuitOpenException("Error pattern detected");

// Get current failure count
int failures = breaker.GetFailureCount(coordinator);

SignalDrivenBackpressure

**Paquete: ** mayormente lúcido.efemeral.patrones.backpresión

Gestión de profundidad de cola con aplazamiento automático en señales de contrapresión.

using Mostlylucid.Ephemeral.Patterns.Backpressure;

var sink = new SignalSink();

await using var coordinator = SignalDrivenBackpressure.Create<WorkItem>(
    async (item, ct) => await ProcessAsync(item, ct),
    sink,
    maxConcurrency: 4);

// Enqueue work
await coordinator.EnqueueAsync(item);

// When downstream is slow
sink.Raise("backpressure.downstream");  // New work auto-defers

// When recovered
sink.Retract("backpressure.downstream"); // Work resumes

ControladoFanOut

**Paquete: ** mayormente lúcido.efemeral.patrones.controladofanout

Gating global + por llave para el paralelismo controlado.

using Mostlylucid.Ephemeral.Patterns.ControlledFanOut;

await using var fanout = new ControlledFanOut<string, Request>(
    keySelector: req => req.TenantId,
    body: async (req, ct) => await ProcessAsync(req, ct),
    maxGlobalConcurrency: 100,  // Total parallel across all tenants
    perKeyConcurrency: 5);      // Max 5 parallel per tenant

// Items for same tenant processed with limit
await fanout.EnqueueAsync(requestA);  // Tenant1
await fanout.EnqueueAsync(requestB);  // Tenant1 - waits if 5 already running
await fanout.EnqueueAsync(requestC);  // Tenant2 - parallel with Tenant1

await fanout.DrainAsync();

AdaptiveRateService

**Paquete: ** principalmente lúcido.efemeral.patrones.adaptiverate

Limitación de velocidad impulsada por señal con retroceso automático.

using Mostlylucid.Ephemeral.Patterns.AdaptiveRate;

await using var service = new AdaptiveRateService<ApiRequest>(
    async (req, ct) => await CallApiAsync(req, ct),
    maxConcurrency: 8);

// Process with automatic rate limit handling
await service.ProcessAsync(request);

// When API returns 429, emit signal with retry-after
// Signal: "rate-limit:500ms"
// Service auto-parses and delays

Console.WriteLine($"Pending: {service.PendingCount}, Active: {service.ActiveCount}");

DinámicaCondiciónDemo

**Paquete: ** En su mayoría lúcido.efemeral.patrones.condición dinámica

Escalado de concurrencia de tiempo de ejecución basado en señales de carga.

using Mostlylucid.Ephemeral.Patterns.DynamicConcurrency;

var sink = new SignalSink();

await using var demo = new DynamicConcurrencyDemo<WorkItem>(
    async (item, ct) => await ProcessAsync(item, ct),
    sink,
    minConcurrency: 2,
    maxConcurrency: 32,
    scaleUpPattern: "load.high",
    scaleDownPattern: "load.low");

await demo.EnqueueAsync(item);

// Concurrency adjusts automatically based on signals
sink.Raise("load.high");  // Concurrency doubles (up to max)
sink.Raise("load.low");   // Concurrency halves (down to min)

Console.WriteLine($"Current concurrency: {demo.CurrentMaxConcurrency}");

await demo.DrainAsync();

KeyedPriorityFanOut

**Paquete: ** Principalmente lucid.efemeral.patterns.keyedpriorityfanout

Calles prioritarias con orden de pedido por llave preservado.

using Mostlylucid.Ephemeral.Patterns.KeyedPriorityFanOut;

await using var fanout = new KeyedPriorityFanOut<string, UserCommand>(
    keySelector: cmd => cmd.UserId,
    body: async (cmd, ct) => await HandleCommand(cmd, ct),
    maxConcurrency: 32,
    perKeyConcurrency: 1,  // Sequential per user
    maxPriorityDepth: 100);

// Normal lane
await fanout.EnqueueAsync(normalCommand);

// Priority lane - jumps the queue for that user
bool accepted = await fanout.EnqueuePriorityAsync(urgentCommand);

// Check lane depths
var counts = fanout.PendingCounts;
Console.WriteLine($"Priority: {counts.Priority}, Normal: {counts.Normal}");

await fanout.DrainAsync();

ReactivaFanOutPipeline

**Paquete: ** principalmente lucid.efemeral.patterns.reactivefanout

Gasoducto de dos etapas con contrapresión automática.

using Mostlylucid.Ephemeral.Patterns.ReactiveFanOut;

await using var pipeline = new ReactiveFanOutPipeline<WorkItem>(
    stage2Work: async (item, ct) => await SlowProcessing(item, ct),
    preStageWork: async (item, ct) => await FastPreprocessing(item, ct),
    stage1MaxConcurrency: 8,
    stage1MinConcurrency: 1,
    stage2MaxConcurrency: 4,
    backpressureThreshold: 32,  // Throttle when stage2 has 32+ pending
    reliefThreshold: 8);        // Resume when stage2 drops below 8

await pipeline.EnqueueAsync(item);

// Stage1 auto-throttles when stage2 backs up
Console.WriteLine($"Stage1 concurrency: {pipeline.Stage1CurrentMaxConcurrency}");
Console.WriteLine($"Stage2 pending: {pipeline.Stage2Pending}");

await pipeline.DrainAsync();

Detector de anomalías de la señal

**Paquete: ** mayormente lúcido.efemeral.patrones.anomalydetector

Detección de anomalías en la ventana en movimiento.

using Mostlylucid.Ephemeral.Patterns.AnomalyDetector;

var sink = new SignalSink();

var detector = new SignalAnomalyDetector(
    sink,
    pattern: "error.*",
    threshold: 5,
    window: TimeSpan.FromSeconds(10));

// Check for anomalies
if (detector.IsAnomalous())
{
    Console.WriteLine("Anomaly detected! Too many errors.");
    TriggerAlert();
}

// Get current match count
int errorCount = detector.GetMatchCount();
Console.WriteLine($"Errors in window: {errorCount}");

SeñalCoordinadaLecturas

**Paquete: ** principalmente lúcido.efemeral.patterns.signalcoordinatedreads

Quiesce lee durante las actualizaciones sin bloqueos duros.

using Mostlylucid.Ephemeral.Patterns.SignalCoordinatedReads;

// Run demo: readers pause when update signal is present
var result = await SignalCoordinatedReads.RunAsync(
    readCount: 10,
    updateCount: 1);

Console.WriteLine($"Reads: {result.ReadsCompleted}, Updates: {result.UpdatesCompleted}");
Console.WriteLine($"Signals: {string.Join(", ", result.Signals)}");

// Manual implementation:
var sink = new SignalSink();

await using var readers = new EphemeralWorkCoordinator<Query>(
    body,
    new EphemeralOptions
    {
        DeferOnSignals = new HashSet<string> { "update.in-progress" },
        Signals = sink
    });

// Readers auto-defer when update is running
sink.Raise("update.in-progress");  // Readers wait
sink.Raise("update.done");         // Readers resume

SeñalizaciónHttpClient

**Paquete: ** mayormentelucid.efemeral.patterns.signalinghttp

Cliente HTTP con señales de progreso.

using Mostlylucid.Ephemeral.Patterns.SignalingHttp;

var httpClient = new HttpClient();
var request = new HttpRequestMessage(HttpMethod.Get, "https://example.com/large-file");

// Create an emitter from your coordinator
// (emitter is any ISignalEmitter - operations implement this)

byte[] data = await SignalingHttpClient.DownloadWithSignalsAsync(
    httpClient,
    request,
    emitter);

// Signals emitted during download:
// - stage.starting
// - progress:0
// - stage.request
// - stage.headers
// - stage.reading
// - progress:25, progress:50, progress:75, progress:100
// - stage.completed

SeñalLogWatter

**Paquete: ** mayormentelucid.efemeral.patterns.signallogwatcher

Mira la ventana de la señal para patrones y dispara callbacks.

using Mostlylucid.Ephemeral.Patterns.SignalLogWatcher;

var sink = new SignalSink();

await using var watcher = new SignalLogWatcher(
    sink,
    onMatch: evt =>
    {
        Console.WriteLine($"Error detected: {evt.Signal} at {evt.Timestamp}");
        AlertOps(evt);
    },
    pattern: "error.*",
    pollInterval: TimeSpan.FromMilliseconds(200));

// Watcher runs in background, calling onMatch for each new error signal
sink.Raise("error.database");    // -> onMatch called
sink.Raise("error.timeout");     // -> onMatch called
sink.Raise("info.started");      // -> ignored (doesn't match pattern)

TelemetríaSeñalHandler

**Paquete: ** principalmente lucid.efemeral.patterns.telemetry

Integración OpenTelemetry/Application Insights.

using Mostlylucid.Ephemeral.Patterns.Telemetry;

// Use in-memory for testing, or implement ITelemetryClient for real telemetry
var telemetry = new InMemoryTelemetryClient();

await using var handler = new TelemetrySignalHandler(telemetry);

// Wire up to coordinator
var options = new EphemeralOptions
{
    OnSignal = signal => handler.OnSignal(signal)
};

// Signals are processed asynchronously
// - "error.*" signals -> TrackExceptionAsync
// - "perf.*" signals -> TrackMetricAsync
// - all signals -> TrackEventAsync

Console.WriteLine($"Queued: {handler.QueuedCount}");
Console.WriteLine($"Processed: {handler.ProcessedCount}");
Console.WriteLine($"Dropped: {handler.DroppedCount}");

// Check recorded events
var events = telemetry.GetEvents();

LongWindowDemo

**Paquete: ** mayormente lúcido.efemeral.patterns.long windowdedemo

Demuestra la configuración de grandes ventanas para pistas de auditoría.

using Mostlylucid.Ephemeral.Patterns.LongWindowDemo;

// Configure coordinator with large tracking window
var options = new EphemeralOptions
{
    MaxTrackedOperations = 10000,
    MaxOperationLifetime = TimeSpan.FromHours(24)
};

SignalReactionShowcase

**Paquete: ** Principalmente lúcido.efemeral.patrones.reacción de la señal mostrar caso

Demuestra patrones de envío de señales y callbacks.

using Mostlylucid.Ephemeral.Patterns.SignalReactionShowcase;

// See source for signal dispatch examples
// Demonstrates OnSignal, OnSignalAsync, CancelOnSignals, DeferOnSignals

Ventana de señal persistente

**Paquete: ** principalmente lucid.efemeral.patterns.persistentwindow

Ventana de señal con persistencia SQLite - sobrevive a los reinicios del proceso.

using Mostlylucid.Ephemeral.Patterns.PersistentWindow;

await using var window = new PersistentSignalWindow(
    "Data Source=signals.db",
    flushInterval: TimeSpan.FromSeconds(30));

// On startup: restore previous signals
await window.LoadFromDiskAsync(maxAge: TimeSpan.FromHours(24));

// Raise signals as normal
window.Raise("order.completed", key: "order-service");
window.Raise("payment.processed", key: "payment-service");

// Query signals
var recentOrders = window.Sense("order.*");

// Signals automatically flush every 30 seconds
// Also flushes on dispose

// Get stats
var stats = window.GetStats();
Console.WriteLine($"In memory: {stats.InMemoryCount}, Flushed: {stats.LastFlushedId}");

Inyección por dependencia

// Register in Startup/Program.cs
services.AddEphemeralWorkCoordinator<WorkItem>(
    async (item, ct) => await ProcessAsync(item, ct),
    new EphemeralOptions { MaxConcurrency = 8 });

// Named coordinators
services.AddEphemeralWorkCoordinator<WorkItem>("priority",
    async (item, ct) => await ProcessPriorityAsync(item, ct));

// Inject and use
public class MyService(IEphemeralCoordinatorFactory<WorkItem> factory)
{
    public async Task DoWork()
    {
        var coordinator = factory.CreateCoordinator();
        await coordinator.EnqueueAsync(new WorkItem());
    }
}

Las raíces DI modernas pueden preferir los ayudantes más cortos tales como services.AddCoordinator<T>(...), services.AddScopedCoordinator<T>(...), o services.AddKeyedCoordinator<T, TKey>(...) ya que leen como normales AddX registraciones; simplemente delegan a los ayudantes específicos de Efímeros bajo el capó.


Marcos de objetivos

  • .NET 6,0, 7,0, 8,0, 9,0, 10.0

Licencia

Unlicencia (dominio público)

Finding related posts...
logo

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