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
Nou dit is mijn obsessie geweest voor de afgelopen week. Zie de vorige delen en wat leidde tot dit; 'Wat als een LRU een uitvoeringscontext was.' Nu is het een set van 30 Nuget pakketten die de meeste belangrijke gelijktijdige uitvoeringspatronen (in Tiny 5-10 lijnpakketten) omvatten. Krijg geweldige adaptieve mogelijkheden met een SIMPLE syntax!
Vind de bron hier: https://github.com/scottgal/mostlylucid.atomen/blob/main/mostlylucid.efemeral/src/mostlylucid.efemeral.complete
Lees het [vorige deel over Signalen ]Hier is de Readme.md van de meest lucid.ephemeral.complete pacakage die zowel de kern meestlucid.ephemeral pakket bevat als alle patronen, 'atomen' (coördinatoren etc.) pakketten in een handige DLL.
OF de kern gebruiken meestal lucid.efemeral een Tiny (letterlijk 10 klassen) die je alle ruwe functionaliteit geeft.
OF als u wilt volledige attribuut gebaseerd async routing met eenvoudige [EphemeralJob] en service.Toevoegen
Dit is waarschijnlijk HET onderwerp van mijn blog vooruit te gaan...je bent gewaarschuwd
All of Mostlylucid.Ephemeral in een single DLL - begrensde async-uitvoering met signaalgebaseerde coördinatie.
dotnet add package mostlylucid.ephemeral.complete
Dit pakket compileert alle kern, atoom, en patroon code in één assemblage. Voor individuele pakketten, zie de links in elke sectie hieronder.
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();
Het vertrouwde services.AddCoordinator<T>() helpers en AddEphemeralSignalJobRunner<T>() Houd service registratie beknopt, laat DI eigenaar van de spoelbak/runner, en maak de nieuwe verantwoordelijkheid/cache/logging verhalen een enkele klik verwijderd.
mostlylucid.ephemeral.complete bundels mostlylucid.ephemeral.attributes, dus attribuut pijpleidingen maken deel uit van de kern
oppervlak. Behandel de loper als een eersteklas signaal consument: gedecoreerde methoden voegen zich bij dezelfde caching, logging, en
pinning verhalen, en elk attribuut kan declareren Priority, job-niveau MaxConcurrency, Lane, Key bronnen, signaal
emissies, pin/expire overrides, en retrieves.
Sleutelattribuutknoppen:
Priority, MaxConcurrency, en Lane werk in deterministische volgorde te houden terwijl hete paden
Blijf uit elkaar.OperationKey, KeyFromSignal, KeyFromPayload, en [KeySource] help u om samen te werken met
betekenisvolle sleutels voor logging, eerlijke planning en diagnostiek.Pin, ExpireAfterMs, AwaitSignals, MaxRetries, en RetryDelayMs laat de verwerkers uitbreiden
hun zichtbaarheid, gate executie totdat afhankelijkheden arriveren, en helen met retrieves terwijl het uitzenden van storingssignalen.EmitOnStart, EmitOnComplete, en EmitOnFailure om downstream stadia te signaleren, log
Wachters, of andere coördinatoren zonder handmatige bedrading.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;
}
}
Deze loper zit nu bij het opstarten en reageert wanneer log.error.* of een uitgezonden signaal raakt de gootsteen. Attribuut
verwerkers kunnen ook sleutels lezen van signalen/payloads, pin werk tot downstream acks, uitstoten voltooiing/mislukking signalen, en
sleuf in rijstroken om te bestellen. Voor DI-eerste setups gebruiken services.AddEphemeralSignalJobRunner<T>() (of de scoped
variant) zodat de loper en spoelbak worden beheerd door de container.
[EfemeralJobs(SignalPrefix = "stage," StandardLane = "pipeline")] openbare verzegelde klasse StageJobs { [EfemeralJob("ingest," EmitOnComplete = new[] { "stage.ingest.done" }) publieke taak IngestAsync(SignalEvent evt) => Console.Out.WriteLineAsync(evt.Signal);
[EphemeralJob("finalize")]
public Task FinalizeAsync(SignalEvent evt) => Console.Out.WriteLineAsync("final stage");
}
var stageSink = nieuwe SignalSink(); wacht met var stageRunner = nieuwe EfemeralSignalJobRunner(stageSink, new[] { nieuwe stagejobs() }); faseSink.Raise("stage.ingest");
Pin-heavy banen kunnen rekenen op ResponsibilitySignalManager.PinUntilQueried (standaard ack patroon responsibility.ack.*) tot en met
hun activiteiten zichtbaar te houden totdat een downstream lezer de lading ophaalt, terwijl OperationEchoMaker/
OperationEchoAtom blijven de laatste signaalstroom, zodat auditors of moleculen nog steeds kunnen proeven de laatste toestand, zelfs na
Het atoom sterft.
mostlylucid.ephemeral.complete bevat ook mostlylucid.ephemeral.atoms.scheduledtasks. Definieer cron of JSON
schema's via ScheduledTaskDefinition (cron, signaal, optioneel key, payload, description, timeZone, format,
runOnStartup, enz.), en laat ScheduledTasksAtom enqueue duurzaam werk door middel van DurableTaskAtom. Elke geplande taak
verhoogt het geconfigureerde signaal in een coordinatorvenster, dus het erft pinning, logging, en verantwoordelijkheid semantiek
terwijl je moleculen of pijpleidingen reageren op de uitgezonden signaalgolf.
Elke DurableTask draagt het schema Name, Signal, facultatief Key, zelfs een getypte Payload, en Description, dus downstream luisteraars weten meteen welke taak liep en welke metadata (bestandsnamen, URL's, enz.) te consumeren. DurableTaskAtom.WaitForIdleAsync() als je gewoon wilt wachten op de huidige uitbarsting van gepland werk te voltooien zonder het atoom te voltooien, houden de scheduler klaar voor de volgende cron tick.
mostlylucid.ephemeral.logging spiegelt Microsoft.Extensions.Inloggen in signalen en vice versa.
SignalLoggerProvider naar uw logger fabriek dus log gebeurtenissen te verhogen log.* signalen, en haak SignalToLoggerAdapter als
U wilt signalen terug in de standaard log pijplijn.
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;
}
}
Gebruik SignalToLoggerAdapter om de resulterende signalen terug te spiegelen in standaard logs zodat uw monitoring stack ziet beide
De zijkanten van de brug.
Pakket: meestal lucid.efemeral
Langlevende werkwachtrij met begrensde concurrency en waarneembaar venster.
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-key sequentiële verwerking - items met dezelfde sleutel verwerkt in volgorde.
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);
Vang resultaten van async operaties.
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);
Meerdere prioritaire rijstroken met configureerbare concurrency per rijstrook.
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
}
Operations zenden signalen uit voor horizontale observeerbaarheid.
// 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
Moeten de resultaten net lang genoeg zichtbaar blijven voor downstreamconsumenten? ResponsibilitySignalManager laat je een speld
operatie totdat er een ack-signaal arriveert (standaard patroon) responsibility.ack.* met sleutel=operationId) Voorzien in een
facultatief description zodat de operatie haar verantwoordelijkheid kan beschrijven, en maxPinDuration naar sierlijk
Zelfverzekerd als de consument nooit komt opdagen.
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 blijft klein (operatie-id, sleutel, signaal, tijdstempel), zodat u kunt opnemen welke minimale toestand u wilt
ongeveer voordat de operatie wordt verzameld.
De coördinator houdt ook een kortstondige echo van de eindsignalen (gesteund via EnableOperationEcho) dat je kunt
inspecteren met GetEchoes() wanneer u de getrimde signaalgolf opnieuw moet afspelen zonder de volledige werking rond te houden.
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 en OperationEchoCapacity laat je in balans brengen hoeveel echo's je houdt en hoe lang ze blijven hangen,
zodat je de laatste woorden kunt herhalen, net lang genoeg om diagnostiek aan de oppervlakte te brengen.
De manager wordt automatisch losgekoppeld als de ack brandt, maar u kunt bellen CompleteResponsibility(operationId) tot beëindiging van de
Verantwoordelijkheid in een vroeg stadium (bijv. bij retrieves). OperationFinalized wanneer het venster ze trimt, dus
abonneer je als je een eindsignaal wilt uitzenden, log diagnostics wilt uitvoeren of de laatste woorden wilt opruimen.
**Pakket: ** meestal lucid.efemeral.atomen.fixedwork
Vaste werker zwembad met statistieken. Minimale API wrapper rond EfemeralWorkCoordinator.
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();
**Pakket: ** meestal lucid.ephemeral.atoms.keyedsequential
Per-key sequentiële verwerking met optionele fair scheduling.
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();
**Pakket: ** meestal lucid.efemeral.atomen.signalaware
Pauzeer of annuleer de opname op basis van omgevingssignalen.
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();
**Pakket: ** meestal lucid.efemeral.atomen.batching
Verzamel items in batches per grootte of tijdsinterval.
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
Exponentiële back-off 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();
Gedeelde configuratie voor opslagatomen (DataStorageConfig, IDataStorageAtom<TKey, TValue>) plus de signaalconventies die bestand, SQLite en PostgreSQL-adapters aansturen.
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");
Gebruik hetzelfde DataStorageConfig met Mostlylucid.Ephemeral.Atoms.Data.Sqlite of Mostlylucid.Ephemeral.Atoms.Data.Postgres implementaties voor duurzame, signaalgestuurde persistentie aangedreven door SQLite/Postgres. Attribuut jobs kunnen zich abonneren op saved.data.{dbname} signalen om downstream werk af te trappen terwijl load.data.{dbname} triggers hydrateren caches.
**Pakket: ** meestal lucid.efemeral.atomen.moleculen
Blueprints gecomponeerd met MoleculeBlueprintBuilder laat u de atomen definiëren (betaling, inventaris, verzending,
kennisgeving) die moet worden uitgevoerd wanneer een signaal zoals order.placed Komt eraan. MoleculeRunner luistert naar de trekker
patroon, creëert een gedeelde MoleculeContext, en voert elke stap uit terwijl u zich abonneert op starten/aanvullen van evenementen.
AtomTrigger wanneer het signaal van een atoom een andere coördinator of molecuul moet starten.
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");
Molecuulstappen kunnen extra signalen oproepen (ctx.Raise("order.shipping.start")) zodat de rest van het systeem de
Baton.
**Pakket: ** meestal lucid.efemeral.atomen.slidingcache
Cache met sliding expiration - toegang tot een resultaat reset zijn 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}");
Pakket: kern (
mostlylucid.ephemeral) Zelfoptimaliserende cache met glijdende TTL op elke hit en uitgebreide TTL voor Hot keys.
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
Tip:
MemoryCachekan worden geconfigureerd voor schuifverval, maar het zendt nooit de warme/koude signalen uit of breidt TTL uit voor hot keys.EphemeralLruCacheis de zelfoptimaliserende standaard in het kernpakket (en inSqliteSingleWriter) wanneer u wilt dat de cache zich focust op de actieve werkende set.
Vang de getypte laatste woorden op die een operatie uitstraalt voordat het wordt getrimd. Het atoom houdt een begrensd venster van signaal
payloads (matching) ActivationSignalPattern / CaptureSignalPattern) en wanneer OperationFinalized vuur dat het produceert
OperationEchoEntry<TPayload> records die u kunt aanhouden via 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");
Attribuut banen alleen verhogen het getypte signaal met welke staat ze kritisch achten, en de maker houdt de werkende set begrensd terwijl je de echo volhoudt.
**Pakket: ** meestal lucid.efemeral.patterns.circuitbreaker
Staatloze stroomonderbreker met signaalgeschiedenisvenster.
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);
**Pakket: ** meestal lucid.efemeral.patterns.backpressure
Wachtrijdieptebeheer met automatisch uitstel op tegendruksignalen.
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
**Pakket: ** meestal lucid.efemeral.patterns.controlled fanout
Global + per-key gating voor gecontroleerd parallelisme.
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();
**Pakket: ** meestal lucid.efemeral.patterns.adaptiverate
Signaalgestuurde snelheidsbeperking met automatische terugslag.
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}");
**Pakket: ** meestal lucid.efemeral.patterns.dynamicconcurrency
Runtime concurrency schalen op basis van belasting signalen.
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();
**Pakket: ** meestal lucid.ephemeral.patterns.keyedpriorityfanout
Prioritaire rijstroken met per-key bestelling bewaard.
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();
**Pakket: ** meestal lucid.efemeral.patterns.reactivefanout
Tweetraps pijpleiding met automatische tegendruk.
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();
**Pakket: ** meestal lucid.efemeral.patterns.anomalydetector
Bewegende-venster anomalie detectie.
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}");
**Pakket: ** meestal lucid.efemeral.patterns.signalcoordinatedreads
Stilte leest tijdens updates zonder harde sloten.
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
**Pakket: ** meestal lucid.efemeral.patterns.signalinghttp
HTTP client met voortgangssignalen.
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
**Pakket: ** meestal lucid.efemeral.patterns.signallogwatcher
Bekijk signaalvenster voor patronen en trigger callbacks.
using Mostlylucid.Ephemeral.Patterns.SignalLogWatcher;
var sink = new SignalSink();
await using var watcher = new SignalLogWatcher(
sink,
onMatch: evt =>
{
Console.WriteLine($"Error detected: {evt.Signal} at {evt.Timestamp}");
AlertOps(evt);
},
pattern: "error.*",
pollInterval: TimeSpan.FromMilliseconds(200));
// Watcher runs in background, calling onMatch for each new error signal
sink.Raise("error.database"); // -> onMatch called
sink.Raise("error.timeout"); // -> onMatch called
sink.Raise("info.started"); // -> ignored (doesn't match pattern)
**Pakket: ** meestal lucid.efemeral.patterns.telemetrie
OpenTelemetrie/Application Insights integratie.
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();
**Pakket: ** meestal lucid.efemeral.patterns.longwindowdemo
Demonstreert grote vensterconfiguratie voor audit trails.
using Mostlylucid.Ephemeral.Patterns.LongWindowDemo;
// Configure coordinator with large tracking window
var options = new EphemeralOptions
{
MaxTrackedOperations = 10000,
MaxOperationLifetime = TimeSpan.FromHours(24)
};
**Pakket: ** meestal lucid.efemeral.patterns.signal reactionshowcase
Demonstreert signaal verzending patronen en terugroepen.
using Mostlylucid.Ephemeral.Patterns.SignalReactionShowcase;
// See source for signal dispatch examples
// Demonstrates OnSignal, OnSignalAsync, CancelOnSignals, DeferOnSignals
**Pakket: ** meestal lucid.efemeral.patterns.persistentwindow
Signaalvenster met SQLite persistentie - overleeft proces herstarten.
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());
}
}
Moderne DI roots kunnen de voorkeur geven aan de kortere helpers zoals services.AddCoordinator<T>(...),
services.AddScopedCoordinator<T>(...), of services.AddKeyedCoordinator<T, TKey>(...) omdat ze als normaal lezen.
AddX registraties; ze delegeren gewoon aan de Efemeral-specifieke helpers onder de kap.
Unlicense (openbaar domein)
© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.