Enimmäkseen lucid.efemeral.fectific; outo samanaikainen järjestelmämalli LRU:n välimuistissa. (Suomi (Finnish))

Enimmäkseen lucid.efemeral.fectific; outo samanaikainen järjestelmämalli LRU:n välimuistissa.

Sunday, 14 December 2025

//

21 minute read

No tämä on ollut pakkomielteeni viimeisen viikon ajan. Katso edelliset osat ja mikä johti tähän; "Mitä jos LRU oli teloituskonteksti." Nyt se on sarja 30 Nuget-pakettia, jotka kattavat tärkeimmät samanaikaiset suoritusmallit (TINY 5-10 -rivin paketeilla). Hanki hämmästyttävät adaptiiviset kyvyt SIMPLE-syntaksilla!

Lähde löytyy täältä: https://github.com/scottgal/mostlylucid.atoms/blob/main/mostlylucid.efemeral/src/mostlylucid.efemeral.täydellinen

Ennakoitsijat

Lue [edellinen osa Signals ]Tässä on Readme.md suurimmaksi osaksi lucid.efemeral.täydellinen pakage, joka sisältää sekä ytimen enimmäkseen lucid.efemeral paketti ja kaikki kuviot, "atoms" (koordinaattorit jne.) paketteja yhdessä kätevässä DLL.

Ydin

OR käyttää ydintä enimmäkseen lusid.efeseral TINY (kirjaimellisesti 10 luokkaa), joka antaa sinulle kaiken raa'an toiminnallisuuden.

Attribuutit ja DI

TAI jos haluat täydellisen attribuuttiin perustuvan async-reitityksen yksinkertaisella [EphemeralJob] ja palvelu.LisääKoordinaattorin tyylirekisteröinti käyttää Enimmäkseen lucid.efemeral.attribuutit.

Tämä on todennäköisesti blogini aihe menossa eteenpäin... sinua on varoitettu

Enimmäkseen lucid.Ephemeral.Täydellinen

NuGet

Kaikki Enimmäkseenlucid.Ephemeral yhdessä DLL - rajattiin asynkkien teloitus signaalipohjaisella koordinoinnilla.

dotnet add package mostlylucid.ephemeral.complete

Tämä paketti kokoaa kaikki ytimen, atomin ja kuvion koodit yhdeksi kokoonpanoksi. kohta alla.


Sisällys


Pikakäynnistys

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

Palveluiden rekisteröinti

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

"Tuttu" services.AddCoordinator<T>() auttajia ja AddEphemeralSignalJobRunner<T>() Pidä palvelurekisteröinti ytimekkäänä, anna DI:n omistaa pesuallas/juoksija ja tee uusista vastuullisuus-/tarkkuuksista yhden klikkauksen päässä.

Ominaisuuslähtöiset työpaikat

mostlylucid.ephemeral.complete nippuja mostlylucid.ephemeral.attributes, joten attribuuttiputkistot ovat osa ydintä Pinta. Kohtele juoksijaa ensiluokkaisena signaalikuluttajana: koristellut menetelmät liittyvät samoihin välilyönteihin, puunkorjuuseen ja Pinning tarinoita, ja jokainen attribuutti voi julistaa Priority, työpaikkataso MaxConcurrency, Lane, Key lähteet, signaali Päästöt, pin/expire ohitetaan ja retroidaan.

Avainominaisuuden nupit:

  • Tilaus ja kaistat: Käyttö Priority, MaxConcurrency, ja Lane pitää työt deterministisessä järjestyksessä, kun on kuumat polut pysyttele erillään.
  • Avainten ja tunnisteiden merkitseminen: OperationKey, KeyFromSignal, KeyFromPayload, ja [KeySource] auttaa ryhmätyössä mielekkäät avaimet kirjautumiseen, reiluun aikataulutukseen ja diagnostiikkaan.
  • Pinning & retries: Pin, ExpireAfterMs, AwaitSignals, MaxRetries, ja RetryDelayMs anna käsittelijöiden pidentää Niiden näkyvyys, porttien toteuttaminen, kunnes riippuvuudet saapuvat, ja paraneminen palautuksilla samalla, kun ne lähettävät vikasignaaleja.
  • Signaalikoreografia: Emit EmitOnStart, EmitOnComplete, ja EmitOnFailure signaloidakseen loppupään vaiheita, loki Tarkkailijat tai muut koordinaattorit ilman manuaalista johdotusta.
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;
    }
}

Juoksija istuu nyt startupissa ja reagoi aina, kun log.error.* tai mikä tahansa signaali osuu altaaseen. Attribuutti Käsittelijät voivat myös lukea avaimia signaaleista/latauksista, nasta työstää alavirtaan asti, lähettää loppuunsaatettuja/viallisia signaaleja ja Lähtö kaistalle tilaamista varten. DI-ensiasennusten käyttöön services.AddEphemeralSignalJobRunner<T>() (tai laajuudeltaan Muunnelma) niin juoksija ja nielu hoidetaan kontilla.

[EphemeralJobs(SignalPrefix = "vaihe", DefaultLane = "putkisto")] julkinen sinetöity luokka StageJobs { [EphemeralJob('ingest', EmitOnComplete = uusi[] { "stage.ingest.done" }] Julkinen tehtävä IngestAsync (SignalEvent evt) = > Console.Out.WriteLineAsync(evt.Signal);

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

}

var stageSink = uusi SignalSink(); odota var stage Runner = uusi EphemeralSignalJobRunner (Sink-vaihe, uusi[] {Uudet StageJobs() }; StageSink.Raise('stage.ingest');

Raskaat työtehtävät voivat luottaa ResponsibilitySignalManager.PinUntilQueried (oletusack-kuvio responsibility.ack.*) Pidä niiden toiminnot näkyvillä, kunnes alavirtalukija hakee hyötykuorman, kun taas OperationEchoMaker/ OperationEchoAtom Jatka lopullista signaalivirtaa, jotta tarkastajat tai molekyylit voivat yhä maistaa viimeistä tilaa vielä sen jälkeenkin. Atomi kuolee.

Aikataulutetut tehtävät

mostlylucid.ephemeral.complete sisältää myös mostlylucid.ephemeral.atoms.scheduledtasks. Määrittele kruunu tai JSON aikataulut kautta ScheduledTaskDefinition (Cron, signaali, valinnainen) key, payload, description, timeZone, format, runOnStartup, jne.), ja anna ScheduledTasksAtom Ennakoi kestävä työ läpi DurableTaskAtom. Jokainen suunniteltu työ nostaa konfiguroidun signaalin koordinaattorin ikkunaan, joten se perii semantiikkaa, kirjautumista ja vastuullisuutta Samalla kun molekyylisi tai määrität putkistosi reagoivat lähettämääsi signaaliaaltoon.

Kaikki DurableTask kantaa aikataulua Name, Signal, valinnainen Key, jopa kirjoitettu Payload, ja Description, joten alajuoksun kuuntelijat tietävät heti, mitä työtä he tekivät ja mitä metatietoja (tiedostonimet, URL-osoitteet jne.) kulutetaan. DurableTaskAtom.WaitForIdleAsync() kun haluat vain odottaa, että nykyinen aikataulutyön puhkeaminen päättyy ilman atomin valmistumista, pitää aikatauluttajan valmiina seuraavaa kroniittipuikkoa varten.

Kirjautuminen ja signaalit

mostlylucid.ephemeral.logging Peilit Microsoft.Extensions.Kirjautuminen signaaleihin ja päinvastoin. Aloita kiinnittämällä SignalLoggerProvider metsuritehtaallesi niin lokitapahtumat nostavat log.* signaaleja, ja koukku SignalToLoggerAdapter jos halutaan, että signaalit virtaavat takaisin normaaliin hirsiputkeen.

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

Käyttö SignalToLoggerAdapter peilata tuloksena olevat signaalit takaisin vakiolokeihin, jotta seurantapinosi näkee molemmat Sillan sivut.


Keskeisiä koordinaattoreita

Paketti: enimmäkseen lusid.efeseral

EphemeralWork -koordinaattori<t>

Pitkäaikainen työjono, jossa on sidottuna valuuttaa ja näkyvä ikkuna.

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

EphemeralKeyedWork Coordinator<Avain, T>

Per-avain peräkkäinen käsittely - kohteet, joissa sama avain käsitellään järjestyksessä.

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

EphemeralResultCoordinator<TInput, Tresult>

Async-toiminnoista saatua tulosta.

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

Ensisijainen työkoordinaattori<t>

Useita ensisijaisia kaistoja, joilla on konfiguroitava vaihtovelkakirja kaistaa kohti.

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

Asetukset (EphemeralOptions)

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
}

Signaalit

Toiminta lähettää signaaleja laaja-alaiseen havainnointiin.

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

Vastuullisuussignaalit ja viimeistely

Täytyykö tuloksia pitää näkyvillä vain riittävän kauan loppupään kuluttajille? ResponsibilitySignalManager Antaa sinun pinnata Toimi kunnes ack-signaali saapuu (oletuskuvio) responsibility.ack.* Avaimella =operationId) Tarjotaan valinnainen description niin operaatio voi kuvailla vastuunsa, ja asettaa maxPinDuration viehkeästi itsestään selvää, jos kuluttaja ei koskaan ilmesty paikalle.

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 pysyy pienenä (toimintatunnus, avain, signaali, aikaleima), jotta voit tallentaa minkä tahansa vähäisen tilan, josta välität noin ennen kuin operaatio kerätään.

Koordinaattori pitää myös lyhyen kaiun lopullisista signaaleista (mahdollisesti EnableOperationEcho) että voit tarkista GetEchoes() kun täytyy toistaa trimmattu signaaliaalto pitämättä koko toimintaa ympärillä.

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 sekä OperationEchoCapacity anna tasapainottaa, kuinka monta kaikua pidät ja kuinka kauan ne viipyvät, Voit siis toistaa viimeiset sanat vain pintadiagnostiikkaan asti.

Johtaja aukeaa automaattisesti, kun ack palaa, mutta voit soittaa CompleteResponsibility(operationId) lopettaakseen Vastuullisuus varhaisessa vaiheessa (esim. rekrytoinnissa). OperationFinalized kun ikkuna leikkaa niitä, niin Tilaa, jos haluat lähettää viimeisen signaalin, kirjautua diagnostiikkaan tai suorittaa viimeisen sanan puhdistuksen.


Atoms (Rakennuslohkot)

FixedWorkAtom

**Paketti: ** Enimmäkseen lucid.efemeral.atoms.kiinteät työt

Kiinteä työntekijäpooli ja tilastot. Minimirajapinta EphemeralWork Coordinatorin ympärillä.

using Mostlylucid.Ephemeral.Atoms.FixedWork;

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

await atom.EnqueueAsync(item);

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

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

// Graceful shutdown
await atom.DrainAsync();

KeyedSequentialAtom

**Paketti: ** enimmäkseen lucid.efemeral.atoms.keyedsequential

Per-avain peräkkäinen käsittely valinnaisella reilulla aikataululla.

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

**Paketti: ** enimmäkseen lucid.efemeral.atoms.signalaware

Keskeytä tai peruuta otto ympäristön signaaleihin perustuen.

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

**Paketti: ** enimmäkseen lucid.efemeral.atoms.atching

Kerää esineet eriin koon tai aikavälin mukaan.

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

Paketti: enimmäkseen lucid.efemeral.atoms.retry

Eksponentiaaliset taustajoukot kokeilevat uudelleen käärepaperia.

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

Tietojen säilytysatomit

Paketti: enimmäkseen lucid.efemeral.atoms.data

Varastoatomien yhteinen kokoonpano (DataStorageConfig, IDataStorageAtom<TKey, TValue>) sekä signaalikäytännöt, jotka ajavat tiedoston, SQLiten ja PostgreSQL:n adapterit.

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

Käytä samaa DataStorageConfig yy) kanssa, kun Mostlylucid.Ephemeral.Atoms.Data.Sqlite tai Mostlylucid.Ephemeral.Atoms.Data.Postgres SQLite/Postgres -ohjelmiston avulla kestävälle, signaalivetoiselle sinnikkyydelle. saved.data.{dbname} signaaleja alkupään töiden käynnistämisestä load.data.{dbname} laukaisee hydraattivarastot.


MoleculeRunner & AtomTrigger

**Paketti: ** enimmäkseen lucid.efemeral.atoms.molecules

Piirustukset, jotka on laadittu MoleculeBlueprintBuilder anna määritellä atomit (maksu, varasto, laivaus, ilmoitus), jonka pitäisi toimia, kun signaali, kuten order.placed saapuu. MoleculeRunner kuuntelee liipaisinta kuvio, luo jaetun MoleculeContext, ja suorittaa jokaisen vaiheen, kun tilaat aloitus-/täyttötapahtumat. AtomTrigger kun yhden atomin signaalin pitäisi käynnistää toinen koordinaattori tai molekyyli.

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

Molekyyrivaiheet voivat tuoda lisäsignaaleja (ctx.Raise("order.shipping.start")) joten muu järjestelmä poimii tahtipuikko.


SlidingCacheAtom

**Paketti: ** enimmäkseen lucid.efemeral.atoms.liukumäkicache

Cache liukuvalla vanhenemisella - tulos nollaa TTL:n.

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

EphemeralLruCache

Paketti: ydin (mostlylucid.ephemeral) – itseoptimoiva välimuisti liukuvalla TTL:llä jokaisen osuman kohdalla ja laajentamalla TTL:ää Kuumat avaimet.

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

Vihje: MemoryCache Voidaan määrittää liukuvan vanhenemisen varalta, mutta se ei koskaan lähetä kuumia/kylmiä signaaleja tai pidentää TTL:ää kuumille avaimille. EphemeralLruCache on itsekeskeinen oletus ydinpaketissa (ja SqliteSingleWriter) aina, kun haluat välimuistin keskittyvän aktiiviseen työsarjaan.

Kaikuluotain

Paketti: enimmäkseen lucid.efemeral.atoms.echo

Napata kirjoitettuja sanoja, jotka toiminto lähettää ennen kuin se leikataan. Atomi pitää rajattu ikkuna signaalin Hyötykuormat (vastaavia) ActivationSignalPattern / CaptureSignalPattern) ja milloin OperationFinalized sen tuottamat tulipalot OperationEchoEntry<TPayload> levyt voit sinnitellä kautta 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");

Attribuutit vain nostavat kirjoitettua signaalia kriittiseksi katsomallaan tavalla, ja tekijä pitää työsettinsä Rajattuna, kun jatkat kaikua.


Mallit (valmiit käyttöön)

SignalBasedCircuitBreaker

**Paketti: ** Enimmäkseen lucid.efemeral.patterns.katkaisija

Valtioton virrankatkaisin signaalihistorian ikkunaa käyttäen.

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

SignalDrivenTakapaine

**Paketti: ** enimmäkseen lucid.efemeral.patterns.backpression

Viivytä syvyyshallintaa automaattisella lykkäyksellä vastapainesignaaleihin.

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

HallittuFanOut

**Paketti: ** enimmäkseen lucid.efemeral.patterns.controlledfanout

Maailmanlaajuinen + per-avain gating kontrolloidulle rinnakkaisuudelle.

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

**Paketti: ** enimmäkseen lucid.efemeral.patterns.adaptiverate

Signaaliohjattu nopeus rajoitti automaattiohjauksella.

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

**Paketti: ** Enimmäkseen lucid.efemeral.patterns.dynaaminen valuutta

Runtime concurrency -skaalaus, joka perustuu kuormitussignaaleihin.

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

AvainpriorityFanout

**Paketti: ** Enimmäkseen lucid.efemeral.patterns.keyedprioriteettifanout

Priority kaistat per-avain tilaamalla säilytetty.

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

**Paketti: ** enimmäkseen lucid.efemeral.patterns.reactivefanout

Kaksivaiheinen putki automaattisella vastapaineella.

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

SignalAnomalyDetector

**Paketti: ** enimmäkseen lucid.efemeral.patterns.anomalydetector

Liikkuvan ikkunan poikkeamahavaitseminen.

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

Signaalien yhteensovitetut lukemat

**Paketti: ** Enimmäkseen lucid.efemeral.patterns.signalcoordinatedreads

Queesce lukee päivityksen aikana ilman kovia lukkoja.

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

SignalingHttpClient

**Paketti: ** enimmäkseen lucid.efemeral.patterns.signalinghttp: / / www.iltasanomat.fi / haku /? search-term = signaling

HTTP-asiakas edistyksellisillä signaaleilla.

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

Signallog Watcher

**Paketti: ** enimmäkseen lucid.efemeral.patterns.signallogwatcher

Katso merkkiikkunaa kuvioiden ja takaisinkutsujen varalta.

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

**Paketti: ** enimmäkseen lucid.efemeral.patterns.telemetria

OpenTelemetria / Application Insights -integraatio.

using Mostlylucid.Ephemeral.Patterns.Telemetry;

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

await using var handler = new TelemetrySignalHandler(telemetry);

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

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

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

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

LongWindowDemo

**Paketti: ** enimmäkseen lucid.efemeral.patterns.longwindowdemo

Osoittaa kirjausketjujen suuren ikkunakokoonpanon.

using Mostlylucid.Ephemeral.Patterns.LongWindowDemo;

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

SignalReactionShowcase

**Paketti: ** Enimmäkseen lucid.efemeral.patterns.signalreaction showcase

Osoittaa signaalien lähetyskuvioita ja takaisinkutsuja.

using Mostlylucid.Ephemeral.Patterns.SignalReactionShowcase;

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

JatkuvaSignalWindow

**Paketti: ** Enimmäkseen lucid.efemeral.patterns.persistantwindow

Signaaliikkuna SQLite sinnikkyydellä - selviää prosessin uudelleenkäynnistämisestä.

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

Riippuvuusruiske

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

Nykyaikaiset DI-juuret saattavat suosia lyhyempiä auttajia, kuten services.AddCoordinator<T>(...), services.AddScopedCoordinator<T>(...), tai services.AddKeyedCoordinator<T, TKey>(...) koska he lukevat normaalisti AddX Ilmoittautumiset; ne vain delegoidaan Ephemeral-erityisille auttajille hupun alla.


Tavoitekehykset

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

Lisenssi

Luvaton (julkinen)

Finding related posts...
logo

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