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
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
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.
OR käyttää ydintä enimmäkseen lusid.efeseral TINY (kirjaimellisesti 10 luokkaa), joka antaa sinulle kaiken raa'an toiminnallisuuden.
TAI jos haluat täydellisen attribuuttiin perustuvan async-reitityksen yksinkertaisella [EphemeralJob] ja palvelu.Lisää
Tämä on todennäköisesti blogini aihe menossa eteenpäin... sinua on varoitettu
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.
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();
"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ä.
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:
Priority, MaxConcurrency, ja Lane pitää työt deterministisessä järjestyksessä, kun on kuumat polut
pysyttele erillään.OperationKey, KeyFromSignal, KeyFromPayload, ja [KeySource] auttaa ryhmätyössä
mielekkäät avaimet kirjautumiseen, reiluun aikataulutukseen ja diagnostiikkaan.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.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.
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.
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.
Paketti: enimmäkseen lusid.efeseral
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();
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);
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);
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");
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
}
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
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.
**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();
**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();
**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();
**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
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();
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.
**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.
**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}");
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:
MemoryCacheVoidaan määrittää liukuvan vanhenemisen varalta, mutta se ei koskaan lähetä kuumia/kylmiä signaaleja tai pidentää TTL:ää kuumille avaimille.EphemeralLruCacheon itsekeskeinen oletus ydinpaketissa (jaSqliteSingleWriter) aina, kun haluat välimuistin keskittyvän aktiiviseen työsarjaan.
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.
**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);
**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
**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();
**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}");
**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();
**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();
**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();
**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}");
**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
**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
**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)
**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();
**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)
};
**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
**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}");
// 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.
Luvaton (julkinen)
© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.