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
Sunday, 14 December 2025
Beh, questa è stata la mia ossessione per la scorsa settimana. Guarda le parti precedenti e cosa ha portato a questo; 'Che cosa se un LRU era un contesto di esecuzione'. Ora è un insieme di 30 pacchetti Nuget che coprono la maggior parte dei principali pattern di esecuzione simultanei (nei pacchetti di linea TINY 5-10). Ottieni capacità adattive sorprendenti con una sintassi SIMPLE!
Per ulteriori informazioni, rivolgersi a:
Leggi il [parte precedente sui segnali ]Per qualche intuizione sui suoi usi. Presentato qui è il Readme.md dal pacakage per lo piùlucid.efemeral.complete che contiene sia il nucleo per lo piùlucid.efemeral pacchetto e tutti i modelli, 'atoms' (coordinatori ecc) pacchetti in un comodo DLL.
O utilizzare il nucleo per lo piùlucid.efetheral un TINY (letteralmente 10 classi) che ti dà tutta la funzionalità raw.
O se si desidera un routing async basato su attributi completi con semplice [EphemeralJob] e servizio.Aggiungi
Questo è probabilmente L'argomento del mio blog che va avanti ... siete stati avvertiti
Tutto di Mostlylucid.Effimero in una singola DLL - esecuzione asincrona limitata con coordinamento basato sul segnale.
dotnet add package mostlylucid.ephemeral.complete
Questo pacchetto compila tutto il codice core, atom e pattern in un unico insieme. Per i singoli pacchetti, vedere i collegamenti in ogni sezione sottostante.
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();
Il familiare services.AddCoordinator<T>() helpers e AddEphemeralSignalJobRunner<T>() mantenere la registrazione del servizio conciso, lasciare DI possedere il lavello / corridore, e rendere le nuove storie responsabilità / cache / logging un solo clic.
mostlylucid.ephemeral.complete fasci mostlylucid.ephemeral.attributes, quindi le condotte di attributo fanno parte del nucleo
superficie. Trattare il corridore come un consumatore di segnale di prima classe: metodi decorati uniscono la stessa cache, registrazione, e
storie di pinning, e ogni attributo può dichiarare Priority, livello di lavoro MaxConcurrency, Lane, Key sorgenti, segnale
emissioni, pin/expire overrides, e ritries.
Pulsanti degli attributi chiave:
Priority, MaxConcurrency, e Lane per mantenere il lavoro in ordine deterministico mentre percorsi caldi
Restate separati.OperationKey, KeyFromSignal, KeyFromPayload, e [KeySource] aiutare a lavorare in gruppo con
chiavi significative per la registrazione, la programmazione equa e la diagnostica.Pin, ExpireAfterMs, AwaitSignals, MaxRetries, e RetryDelayMs consentire ai gestori di estendere
la loro visibilità, l'esecuzione del cancello fino all'arrivo delle dipendenze, e guarire con i tentativi durante l'emissione di segnali di guasto.EmitOnStart, EmitOnComplete, e EmitOnFailure per segnalare le fasi a valle, log
Osservatori, o altri coordinatori senza cablaggio manuale.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;
}
}
Questo corridore ora siede all'avvio e reagisce ogni volta che log.error.* o qualsiasi segnale emesso colpisce il lavandino. Attributo
i gestori possono anche leggere le chiavi da segnali/payloads, lavorare pin fino a downstream acks, emettere segnali di completamento/fallimento, e
slot in corsie per l'ordine. Per l'uso di impostazioni DI-first services.AddEphemeralSignalJobRunner<T>() (o il campo d'applicazione)
variante) in modo che il corridore e lavello sono gestiti dal contenitore.
[EphemeralJobs(SignalPrefix = "stage," DefaultLane = "pipeline")] pubblico sigillato classe StageLavori { [EfemeralJob("ingest," EmitOnComplete = new[] { "stage.ingest.done" }] attività pubblica IngestAsync(SignalEvent evt) => Console.Out.WriteLineAsync(evt.Signal);
[EphemeralJob("finalize")]
public Task FinalizeAsync(SignalEvent evt) => Console.Out.WriteLineAsync("final stage");
}
var stageSink = nuovo SignalSink(); attendere utilizzando var stageRunner = nuovo EphemeralSignalJobRunner(stageSink, nuovo[] {new StageJobs() }); stadioSink.Raise("stage.ingest");
I posti di lavoro più pesanti possono contare su ResponsibilitySignalManager.PinUntilQueried (default ack pattern) responsibility.ack.*) a
mantenere le loro operazioni visibili fino a quando un lettore a valle prende il carico utile, mentre OperationEchoMaker/
OperationEchoAtom perseverare nel flusso di segnale finale in modo che i revisori o le molecole possano ancora assaporare l'ultimo stato anche dopo
l'atomo muore.
mostlylucid.ephemeral.complete contiene anche mostlylucid.ephemeral.atoms.scheduledtasks. Definisci cron o JSON
orari attraverso ScheduledTaskDefinition (crono, segnale, optional) key, payload, description, timeZone, format,
runOnStartup, ecc.), e lasciare ScheduledTasksAtom enqueue lavoro durevole attraverso DurableTaskAtom. Ogni lavoro programmato
solleva il segnale configurato all'interno di una finestra di coordinamento, quindi eredita pinning, logging e responsabilità semantica
mentre le molecole o le condotte di attributo rispondono all'onda di segnale emessa.
Ogni DurableTask porta il programma Name, Signal, facoltativo Key, anche un typed Payload, e Description, così gli ascoltatori a valle sanno immediatamente quale lavoro ha funzionato e quali metadati (nomi di file, URL, ecc.) da consumare. DurableTaskAtom.WaitForIdleAsync() quando si desidera solo aspettare che l'attuale scoppio di lavoro programmato per finire senza completare l'atomo, mantenendo il programmatore pronto per il prossimo cron tick.
mostlylucid.ephemeral.logging Microsoft.Estensioni.Accedere ai segnali e viceversa. Iniziare collegando
SignalLoggerProvider alla tua fabbrica di logger in modo da registrare gli eventi alzare log.* segnali e gancio SignalToLoggerAdapter se
Vuoi che i segnali tornino nella pipeline di log standard.
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 per riflettere i segnali risultanti di nuovo nei registri standard in modo che il vostro stack di monitoraggio vede entrambi
Lati del ponte.
Pacchetto: per lo piùlucid.efetheral
Coda di lavoro a lunga durata con conto corrente limitato e finestra osservabile.
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();
Elaborazione sequenziale per chiave - elementi con la stessa chiave elaborati in ordine.
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);
Cattura dei risultati da operazioni asincrone.
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);
Piu' corsie prioritarie con valuta configurabile per corsia.
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
}
Le operazioni emettono segnali di osservabilità trasversale.
// 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
Necessità di mantenere i risultati visibili appena abbastanza a lungo per i consumatori a valle? ResponsibilitySignalManager ti permette di appuntare un
operazione fino all'arrivo di un segnale ack (modello predefinito responsibility.ack.* con chiave=operationId). Fornire un
facoltativo description in modo che l'operazione possa descrivere la sua responsabilità, e impostare maxPinDuration a graziosamente
auto-chiaro se il consumatore non si presenta mai.
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 rimane piccolo (operazione id, chiave, segnale, timestamp), in modo da poter registrare qualsiasi stato minimo che ti interessi
circa prima che l'operazione sia raccolta.
Il coordinatore mantiene anche un'eco di breve durata dei segnali finali (abilitato tramite EnableOperationEcho) che puoi
ispezionare con GetEchoes() quando è necessario riprodurre l'onda del segnale tagliato senza mantenere l'intero funzionamento intorno.
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 e OperationEchoCapacity ti lasci bilanciare quanti echi tieni e quanto durano,
in modo da poter riprodurre le parole di ultima parola solo abbastanza a lungo per la diagnostica di superficie.
Il gestore apre automaticamente quando l'ack spara, ma puoi chiamare CompleteResponsibility(operationId) per porre fine alla
Le operazioni sono ancora in aumento. OperationFinalized quando la finestra li taglia, così
abbonarsi se si desidera emettere un segnale finale, diagnostica di registro, o eseguire la pulizia di parole di ultima parola.
**Pacchetto: ** Per lo più Lucid.efetheral.atoms.fixedwork
Piscina operaia fissa con statistiche. API minimale avvolgente intorno EfemeralWorkCoordinatore.
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();
**Pacchetto: ** Per lo piùlucid.efetheral.atoms.keyedsequential
Elaborazione sequenziale per-key con programmazione equa opzionale.
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();
**Pacchetto: ** per lo piùlucid.efemeral.atoms.signalaware
Mettere in pausa o annullare l'assunzione in base ai segnali ambientali.
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();
**Pacchetto: ** per lo piùlucid.efetheral.atoms.batching
Raccogliere gli oggetti in lotti per dimensione o intervallo di tempo.
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
Pacchetto: Per lo piùlucid.efetheral.atoms.retry
Esponential backoff retry wrapper.
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();
Pacchetto: per lo più Lucid.efetheral.atoms.data
Configurazione condivisa per atomi di storage (DataStorageConfig, IDataStorageAtom<TKey, TValue>) più le convenzioni di segnale che guidano file, SQLite, e adattatori 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");
Usa lo stesso DataStorageConfig con Mostlylucid.Ephemeral.Atoms.Data.Sqlite oppure Mostlylucid.Ephemeral.Atoms.Data.Postgres implementazioni per la persistenza durevole, basata sul segnale alimentato da SQLite/Postgres. Attribute jobs can subscribe to saved.data.{dbname} segnali per iniziare il lavoro a valle mentre load.data.{dbname} innesca cache idratate.
**Pacchetto: ** per lo piùlucid.efemeral.atoms.molecole
Bluprint composti con MoleculeBlueprintBuilder permette di definire gli atomi (pagamento, inventario, spedizione,
notifica) che dovrebbe essere eseguito quando un segnale come order.placed Arriva. MoleculeRunner ascolta per il grilletto
pattern, crea una condivisione MoleculeContext, ed esegue ogni passo mentre ti iscrivi per avviare/completare eventi. Usa
AtomTrigger quando il segnale di un atomo dovrebbe avviare un altro coordinatore o molecola.
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");
Le fasi Molecule possono sollevare segnali aggiuntivi (ctx.Raise("order.shipping.start")) in modo che il resto del sistema raccoglie il
Manganello.
**Pacchetto: ** per lo piùlucid.efetheral.atoms.slidingcache
Cache con scadenza scorrevole - l'accesso a un risultato ripristina il suo 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}");
Pacchetto: nucleo (
mostlylucid.ephemeral) cache auto-ottimizzante con TTL scorrevole su ogni colpo e TTL esteso per Chiavi calde.
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
Suggerimento:
MemoryCachepuò essere configurato per scadenza scorrevole, ma non emette mai i segnali caldo/freddo o estende TTL per le chiavi calde.EphemeralLruCacheè l'autoottimizzazione predefinita nel pacchetto core (e inSqliteSingleWriter) ogni volta che vuoi che la cache si concentri sul set di lavoro attivo.
Pacchetto: per lo piùlucid.efetheral.atoms.echo
Catturare le parole dattilografate che un'operazione emette prima di essere rifilata. L'atomo mantiene una finestra di segnale delimitata
payload (matching) ActivationSignalPattern / CaptureSignalPattern) e quando OperationFinalized incendi che produce
OperationEchoEntry<TPayload> record che puoi persistere tramite 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");
I lavori di Attributo semplicemente alzano il segnale digitato con qualsiasi stato ritengano critico, e il creatore mantiene il set di lavoro limitato mentre persistete l'eco.
**Pacchetto: ** per lo piùlucid.efemeral.patterns.circuitbreaker
Interruttore senza stato che usa la finestra della cronologia dei segnali.
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);
**Pacchetto: ** per lo piùlucid.efemeral.patterns.backpressure
Gestione della profondità di coda con rinvio automatico sui segnali di contropressione.
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
**Pacchetto: ** Per lo più lucid.efemeral.patterns.fanout controllato
Globale + gating per chiave per parallelismo controllato.
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();
**Pacchetto: ** per lo piùlucid.efemeral.patterns.adaptiverate
Limitazione della velocità a segnale con backoff automatico.
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}");
**Pacchetto: ** Per lo piùlucid.efemeral.patterns.dynamicconcurrency
Runtime concurrency scalata in base ai segnali di carico.
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();
**Pacchetto: ** Per lo piùlucid.efemeral.patterns.keyedpriorityfanout
Linee di priorità con ordine per chiave conservata.
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();
**Pacchetto: ** per lo piùlucid.efemeral.patterns.reactfanout
Conduttura a due stadi con contropressione automatica.
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();
**Pacchetto: ** Per lo piùlucid.efemeral.patterns.anomalydetector
Rilevamento delle anomalie della finestra mobile.
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}");
**Pacchetto: ** Per lo piùlucid.efetheral.patterns.signalreadscoordinati
Quiesce legge durante gli aggiornamenti senza serrature rigide.
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
**Pacchetto: ** Perlopiùlucid.efemeral.patterns.signalinghttp
Client HTTP con segnali di avanzamento.
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
**Pacchetto: ** Per lo piùlucid.efemeral.patterns.signallogwatcher
Finestra del segnale di osservazione per i modelli e le richiamate di attivazione.
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)
**Pacchetto: ** per lo piùlucid.efemeral.patterns.telemetry
OpenTelemetry/Application Insights integration.
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();
**Pacchetto: ** per lo piùlucid.efemeral.patterns.longwindowdemo
Dimostra la grande configurazione delle finestre per i percorsi di audit.
using Mostlylucid.Ephemeral.Patterns.LongWindowDemo;
// Configure coordinator with large tracking window
var options = new EphemeralOptions
{
MaxTrackedOperations = 10000,
MaxOperationLifetime = TimeSpan.FromHours(24)
};
**Pacchetto: ** Per lo più lucido.efemero.patterns.signalereazioni vetrinacaso
Dimostra gli schemi di invio del segnale e i richiami.
using Mostlylucid.Ephemeral.Patterns.SignalReactionShowcase;
// See source for signal dispatch examples
// Demonstrates OnSignal, OnSignalAsync, CancelOnSignals, DeferOnSignals
**Pacchetto: ** per lo più Lucid.efetheral.patterns.persistentwindow
Finestra del segnale con persistenza SQLite - sopravvive ai riavvii di processo.
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());
}
}
Le radici moderne dell'DI possono preferire gli helper più corti come services.AddCoordinator<T>(...),
services.AddScopedCoordinator<T>(...), oppure services.AddKeyedCoordinator<T, TKey>(...) dato che leggono come se fossero normali.
AddX le registrazioni; semplicemente delegano agli aiutanti specifici dell'effimero sotto il cofano.
Unlicense (pubblico dominio)
© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.