Det här har varit min besatthet under den senaste veckan. Se de tidigare delarna och vad som ledde till detta; "Tänk om en LRU var en exekveringskontext". Nu är det en uppsättning av 30 Nuget-paket som täcker de flesta större samtidiga exekveringsmönster (i TINY 5-10 linjepaket). Få fantastiska adaptiva funktioner med en SIMPLE syntax!
Hitta källan här: https://github.com/scottgal/mostlylucid.atoms/blob/main/mostlylucid.ephemeral/src/mostlylucid.ephemeral.complete
Läs [Föregående del om signaler ]för viss insikt i dess användningsområden. Presenteras här är Readme.md från de mestadels lucid.ephemeral.fullständig pacakage som innehåller både kärnan mestadelslucid.ephemeral paket och alla mönster, "atomer" (koordinatorer etc) paket i en bekväm DLL.
ELLER använd kärnan mestadels lucid.ephemeralt en TINY (bokstavligen 10 klasser) som ger dig alla rå funktionalitet.
ELLER om du vill ha full attribut baserad async routing med enkel [EphemeralJob] och service.Lägg till
Detta är sannolikt ämnet för min blogg framöver ... du har blivit varnad på
Hela Mostlylucid.Ephemeral i en enda DLL - begränsad async-utförande med signalbaserad koordination.
dotnet add package mostlylucid.ephemeral.complete
Paketet sammanställer alla kärn-, atom- och mönsterkoder i en enhet. För enskilda paket, se länkarna i varje avsnitt nedan.
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();
Den välbekanta services.AddCoordinator<T>() hjälpare och AddEphemeralSignalJobRunner<T>() hålla registreringen av tjänsten koncis, låta DI äga diskbänken/runner, och göra det nya ansvaret / cache / loggning berättelser ett enda klick bort.
mostlylucid.ephemeral.complete buntar mostlylucid.ephemeral.attributes, så attribut rörledningar är en del av kärnan
Yta. Behandla löparen som en förstklassig signal konsument: dekorerade metoder ansluter samma caching, loggning, och
pinning berättelser, och varje attribut kan deklarera Priority, arbetsnivå MaxConcurrency, Lane, Key källor, signal
utsläpp, pin/expirera åsidosättanden och retries.
Nyckelattributrattar:
Priority, MaxConcurrency, och Lane att hålla arbetet i deterministisk ordning medan heta vägar
Håll dig ifrån varandra.OperationKey, KeyFromSignal, KeyFromPayload, och [KeySource] hjälper dig att arbeta tillsammans med
meningsfulla nycklar för loggning, rättvis schemaläggning och diagnostik.Pin, ExpireAfterMs, AwaitSignals, MaxRetries, och RetryDelayMs Låt handhavarna förlänga
Deras sikt, grindutförande tills beroenden anländer och läker med retries medan de avger felsignaler.EmitOnStart, EmitOnComplete, och EmitOnFailure för att signalera nedströmsfaser, logg
Bevakare eller andra koordinatorer utan manuell ledning.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;
}
}
Den här löparen sitter nu vid start och reagerar när som helst log.error.* eller någon avgiven signal träffar diskhon. Attribut
hanterare kan också läsa tangenter från signaler/payloads, pin arbete till nedströms acks, avger slutförande/underkända signaler, och
plats i köer för beställning. För DI-första inställningar användning services.AddEphemeralSignalJobRunner<T>() (eller tillämpningsområdet
variant) så löpare och handfat hanteras av behållaren.
[EfemeralJobs(SignalPrefix = "steg", StandardLane = "pipeline")] offentlig förseglad klass StageJobs { [EfemeralJob("ingest", EmitOnComplete = nytt[] {"stage.ingest.done"}]] Public Task IngestAsync(SignalEvent evt) => Console.Out.WriteLineAsync(evt.Signal);
[EphemeralJob("finalize")]
public Task FinalizeAsync(SignalEvent evt) => Console.Out.WriteLineAsync("final stage");
}
var stageSink = ny signalSink(); väntar med var stageRunner = nya EfemeralSignalJobRunner(stageSink, ny[] { nya StageJobs()}); ScenSink.Raise("stage.ingest");
Pin-tunga jobb kan lita på ResponsibilitySignalManager.PinUntilQueried (förvalt ack- mönster) responsibility.ack.*) till
hålla sin verksamhet synlig tills en nedströmsläsare hämtar nyttolasten, medan OperationEchoMaker/
OperationEchoAtom envisas den slutliga signalströmmen så revisorer eller molekyler kan fortfarande smaka det sista tillståndet även efter
Atomen dör.
mostlylucid.ephemeral.complete innehåller även mostlylucid.ephemeral.atoms.scheduledtasks. Definiera cron eller JSON
scheman via ScheduledTaskDefinition (krona, signal, valfritt key, payload, description, timeZone, format,
runOnStartup, etc.), och låt ScheduledTasksAtom ett hållbart arbete genom DurableTaskAtom. Varje schemalagt jobb
höjer den inställda signalen inne i ett koordinatorfönster, så det ärver pinning, loggning och ansvar semantik
medan dina molekyler eller attribut pipelines reagerar på den utsända signalvågen.
Var och en DurableTask bär schemat Name, Signal, frivillig Key, även en typad Payload, och Description, så nedströms lyssnare omedelbart vet vilket jobb som kördes och vilka metadata (filnamn, webbadresser, etc.) att konsumera. DurableTaskAtom.WaitForIdleAsync() när du bara vill vänta på att den aktuella explosionen av schemalagt arbete att avsluta utan att slutföra atomen, hålla schemaläggaren redo för nästa cron fästing.
mostlylucid.ephemeral.logging speglar Microsoft.Extensions.Logga in i signaler och vice versa. Börja med att bifoga
SignalLoggerProvider till din logger fabrik så logg händelser höja log.* signaler och krok SignalToLoggerAdapter om
Du vill att signalerna ska flyta tillbaka in i den vanliga log pipelinen.
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;
}
}
Användning SignalToLoggerAdapter för att spegla de resulterande signalerna tillbaka till standardloggar så att din övervakning stack ser båda
sidor av bron.
Förpackning: mestadels lucid.ephemeralt
Långlivad arbetskö med avgränsad konvergens och observerbart fönster.
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-nyckel sekventiell behandling - objekt med samma nyckel behandlas i ordning.
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);
Fånga resultat från async-operationer.
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);
Flera prioriterade körfält med konfigurerbar konvergens per körfält.
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
}
Verksamheten avger signaler för tvärgående observerbarhet.
// 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
Behöver resultaten vara tillräckligt synliga för konsumenter i senare led? ResponsibilitySignalManager Låter dig fästa en
drift tills en acksignal anländer (standardmönster) responsibility.ack.* med nyckel=operationId) . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . .
valfritt description så att operationen kan beskriva sitt ansvar, och ställa in maxPinDuration till graciöst
Självklar om konsumenten aldrig dyker upp.
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 förblir liten (operation id, nyckel, signal, tidsstämpel), så att du kan spela in alla minimala tillstånd du bryr dig
ungefär innan operationen samlas in.
Samordnaren håller också ett kortlivat eko av de slutliga signalerna (aktiverat via EnableOperationEcho) att du kan
inspektera med GetEchoes() när du behöver spela upp den trimmade signalvågen utan att hålla hela operationen runt.
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 och OperationEchoCapacity Låt dig balansera hur många ekon du håller och hur länge de dröjer,
så att du kan spela upp de sista orden bara tillräckligt länge för att ytdiagnostik.
Chefen lossar automatiskt när acken avfyras, men du kan ringa CompleteResponsibility(operationId) för att avsluta
ansvar tidigt (t.ex., på retries). Operationer fortfarande höja OperationFinalized när fönstret trimmar dem, så
prenumerera om du vill avge en slutlig signal, logg diagnostik, eller köra på sista ord på rensning.
**Förpackning: ** Mest lucid.ephemeral.atoms.fixed arbete
Fast arbetskraft pool med statistik. Minimal API omslag runt 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();
**Förpackning: ** mestadels lucid.ephemeral.atoms.keyedsequential
Per nyckel sekventiell behandling med valfri rättvis schemaläggning.
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();
**Förpackning: ** mestadelslucid.ephemeral.atoms.signalaware
Pausa eller avbryta intaget baserat på omgivande signaler.
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();
**Förpackning: ** mestadels lucid.ephemeral.atomer.batching
Samla föremål i omgångar efter storlek eller tidsintervall.
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
Förpackning: mestadels lucid.ephemeral.atoms.retry
Exponentiellt bakslagspapper.
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();
Förpackning: Mest lucid.ephemeral.atomer.data
Delad konfiguration för lagringsatomer (DataStorageConfig, IDataStorageAtom<TKey, TValue>) plus de signalkonventioner som driver fil, SQLite, och PostgreSQL-adaptrar.
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");
Använd samma DataStorageConfig där Mostlylucid.Ephemeral.Atoms.Data.Sqlite eller Mostlylucid.Ephemeral.Atoms.Data.Postgres implementationer för hållbar, signaldriven uthållighet drivs av SQLite/Postgres. Attribut jobb kan prenumerera på saved.data.{dbname} signaler för att starta nedströms arbete medan load.data.{dbname} utlöser hydrerade cacheminnen.
**Förpackning: ** mestadels lucid.ephemerala.atomer.molekyler
Planritningar bestående av MoleculeBlueprintBuilder låt dig definiera atomerna (betalning, inventering, frakt,
anmälan) som bör köras när en signal som order.placed Kommer hit. MoleculeRunner lyssnar efter avtryckaren
mönster, skapar en delad MoleculeContext, och utför varje steg medan du prenumererar på start / komplettering händelser. Använd
AtomTrigger När en atoms signal ska starta en annan koordinator eller molekyl.
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");
Molekylsteg kan höja ytterligare signaler (ctx.Raise("order.shipping.start")) så resten av systemet plockar upp
Baton, tack.
**Förpackning: ** mestadels lucid.ephemeral.atomer.glidande cache
Cache med glidande utgång - få tillgång till ett resultat återställer sin 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}");
Förpackning: Kärna (
mostlylucid.ephemeral) – Självoptimerande cache med glidande TTL på varje träff och förlängd TTL för Snygga nycklar.
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
Tips:
MemoryCachekan konfigureras för att glida ut, men den avger aldrig de varma/kalla signalerna eller förlänger TTL För heta nycklar.EphemeralLruCacheär den självoptimerande standard i kärnpaketet (och iSqliteSingleWriter) När du vill att cachen ska fokusera på den aktiva arbetsuppsättningen.
Förpackning: mestadels lucid.ephemeral.atoms.echo
Fånga den typade "sista orden" som en operation avger innan den klipps. Atomen håller ett avgränsat fönster av signal
nyttolast (matchning) ActivationSignalPattern / CaptureSignalPattern) och när OperationFinalized eldar den producerar
OperationEchoEntry<TPayload> register du kan hålla i 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");
Attribut jobb bara höja den typade signalen med vilket tillstånd de anser kritiskt, och tillverkaren håller arbetssättet avgränsad medan du envisas med ekot.
**Förpackning: ** mestadels lucid.ephemeral.patterns.kretsbrytare
Statslös strömbrytare med hjälp av signalhistorikfönster.
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);
**Förpackning: ** mestadels lucid.ephemeral.patterns.backpress
Köa djuphantering med automatisk uppskov på mottryckssignaler.
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
**Förpackning: ** mestadels lucid.ephemeral.patterns.kontrollerad fanout
Global + per nyckel för kontrollerad parallellism.
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();
**Förpackning: ** mestadels lucid.ephemeral.patterns.adaptiverat
Signalstyrd hastighetsbegränsning med automatisk backoff.
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}");
**Förpackning: ** För det mestalucid.ephemeral.patterns.dynamicconcurrency
Tidskonvergensskalning baserad på belastningssignaler.
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();
**Förpackning: ** mestadels lucid.ephemeral.patterns.keyedpriority fanout
Prioritetsfilar med per-nyckel beställning bevarad.
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();
**Förpackning: ** mestadels lucid.ephemeral.patterns.reactivefanout
Tvåstegsrörledning med automatiskt mottryck.
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();
**Förpackning: ** mestadels lucid.ephemeral.patterns.anomalydetektor
Anomalidetektering av rörliga fönster.
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}");
**Förpackning: ** mestadels lucid.ephemeral.patterns.signalcoordinatedreads
Quiesce läser under uppdateringar utan hårda lås.
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
**Förpackning: ** mest lucid.ephemeral.patterns.signalinghttp
HTTP-klient med framstegssignaler.
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
**Förpackning: ** mestadels lucid.ephemeral.patterns.signallogwatcher
Titta på signalfönster för mönster och utlösa 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)
**Förpackning: ** mestadels lucid.ephemeral.patterns.telemetri
Integrering av OpenTelemetri/Application Insights.
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();
**Förpackning: ** mestadels lucid.ephemeral.patterns.longwindowdemo
Visar stora fönsterkonfigurationer för verifieringsspår.
using Mostlylucid.Ephemeral.Patterns.LongWindowDemo;
// Configure coordinator with large tracking window
var options = new EphemeralOptions
{
MaxTrackedOperations = 10000,
MaxOperationLifetime = TimeSpan.FromHours(24)
};
**Förpackning: ** oftastlucid.ephemeral.patterns.signalreactionshow fall
Demonstrerar signalerar sändningsmönster och samtal.
using Mostlylucid.Ephemeral.Patterns.SignalReactionShowcase;
// See source for signal dispatch examples
// Demonstrates OnSignal, OnSignalAsync, CancelOnSignals, DeferOnSignals
**Förpackning: ** mestadels lucid.ephemeral.patterns.persistentwindow
Signalfönster med SQLite persistens - överlever processomstarter.
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());
}
}
Moderna DI rötter kan föredra de kortare hjälpare som services.AddCoordinator<T>(...),
services.AddScopedCoordinator<T>(...), eller services.AddKeyedCoordinator<T, TKey>(...) eftersom de läser som vanligt
AddX Registreringar; de delegerar helt enkelt till Efemeral-specifika hjälpare under huven.
Olicensierad (offentlig förvaltning)
© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.