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
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.
O utilice el núcleo mayormente lúcido.efemeral un TINY (literalmente 10 clases) que le da toda la funcionalidad en bruto.
O si desea un enrutamiento async basado en atributo completo con sencillo [EphemeralJob] y el servicio.
Esto es probablemente el tema de mi blog que va hacia adelante ... usted ha sido advertido
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.
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 });
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.
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:
Priority, MaxConcurrency, y Lane para mantener el trabajo en orden determinista mientras los caminos calientes
Mantente separado.OperationKey, KeyFromSignal, KeyFromPayload, y [KeySource] ayudarle a trabajar en grupo con
claves significativas para el registro, la programación justa, y diagnósticos.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.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.
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.
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.
Paquete: mayormente lúcido.efemeral
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();
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);
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);
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");
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
}
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
¿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.
**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();
**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();
**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();
**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
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();
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.
**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.
**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}");
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:
MemoryCachese puede configurar para la expiración deslizante, pero nunca emite las señales de calor/frío o extiende TTL para las llaves calientes.EphemeralLruCachees el valor predeterminado de auto-optimización en el paquete principal (y enSqliteSingleWriter) siempre que desee que la caché se centre en el conjunto de trabajo activo.
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.
**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);
**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
**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();
**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}");
**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();
**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();
**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();
**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}");
**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
**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
**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)
**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();
**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)
};
**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
**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}");
// 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ó.
Unlicencia (dominio público)
© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.