Per lo piùlucid.efemeral.complete; uno strano sistema di sistemi paralleli in una cache LRU. (Italiano (Italian))

Per lo piùlucid.efemeral.complete; uno strano sistema di sistemi paralleli in una cache LRU.

Sunday, 14 December 2025

//

25 minute read

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:

Precedenti

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.

Centrale

O utilizzare il nucleo per lo piùlucid.efetheral un TINY (letteralmente 10 classi) che ti dà tutta la funzionalità raw.

Attributi & DI

O se si desidera un routing async basato su attributi completi con semplice [EphemeralJob] e servizio.AggiungiLa registrazione in stile Coordinatore si avvale della pacchetto per lo piùlucid.efemeral.attributes.

Questo è probabilmente L'argomento del mio blog che va avanti ... siete stati avvertiti

Per lo più lucido.Effimero.Completo

NuGetCity name (optional, probably does not need a translation)

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.


Indice


Avvio rapido

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

Registrazione del servizio

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.

Posti di lavoro orientati agli attributi

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:

  • Ordinare & corsie: Uso Priority, MaxConcurrency, e Lane per mantenere il lavoro in ordine deterministico mentre percorsi caldi Restate separati.
  • Tag & tag: OperationKey, KeyFromSignal, KeyFromPayload, e [KeySource] aiutare a lavorare in gruppo con chiavi significative per la registrazione, la programmazione equa e la diagnostica.
  • Pinning & retries: 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.
  • Coreografia del segnale: Emit 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.

Attività pianificate

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.

Registrazione dei & segnali

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.


Coordinatori principali

Pacchetto: per lo piùlucid.efetheral

Coordinatore EphemeralWork<T>

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

EfemeralKeyedWorkCoordinatore<TKey, T>

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

Coordinatore EphemeralResultCoordinator<TInput, TResult>

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

Coordinatore prioritario di lavoro<T>

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

Configurazione (Opzioni effimere)

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
}

Segnali

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

Segnali di responsabilità e finalizzazione

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.


Atomi (blocchi di costruzione)

Atom di lavoro fisso

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

KeyedSequentialAtomCity name (optional, probably does not need a translation)

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

SignalAwareAtom

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

BatchingAtom

**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

RetryAtom

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

Atomi di memorizzazione dati

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.


MolecoleRunner & AtomTrigger

**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.


SlidingCacheAtom

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

EffimeraLruCache

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: MemoryCache può 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 in SqliteSingleWriter) ogni volta che vuoi che la cache si concentri sul set di lavoro attivo.

Echo MakerCity name (optional, probably does not need a translation)

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.


Modelli (pronti all'uso)

SignalBasedCircuitBreaker

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

SegnaleDrivenBackpressure

**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

ControlledFanOut

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

AdaptiveRateService

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

DynamicConcurrencyDemo

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

KeyedPriorityFanOut

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

ReactiveFanOutPipeline

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

Rilevatore di anomalie del segnale

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

SignalCoordinatedReads

**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

SegnalazioneHttpClient

**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

SignalLogWatcher

**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)

TelemetriaSignalHandler

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

LongWindowDemoCity name (optional, probably does not need a translation)

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

Mostracase della reazione del segnale

**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

Finestra di segnale persistente

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

Iniezione di dipendenza

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


Quadri di obiettivi

  • .NET 6.0, 7.0, 8.0, 9.0, 10.0

Licenza

Unlicense (pubblico dominio)

Finding related posts...
logo

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