ως επί το πλείστον διαυγής. epemeral.complete? ένα παράξενο παράλληλο σύστημα μοτίβο συστημάτων σε μια μνήμη LRU. (ελληνικά (Greek))

ως επί το πλείστον διαυγής. epemeral.complete? ένα παράξενο παράλληλο σύστημα μοτίβο συστημάτων σε μια μνήμη LRU.

Sunday, 14 December 2025

//

23 minute read

Λοιπόν αυτή ήταν η εμμονή μου για την περασμένη εβδομάδα. Δείτε τα προηγούμενα μέρη και τι οδήγησε σε αυτό; "Τι θα συμβεί αν μια LRU ήταν ένα πλαίσιο εκτέλεσης." Τώρα είναι ένα σύνολο 30 πακέτα Nuget που καλύπτουν τα πιο σημαντικά πρότυπα ταυτόχρονης εκτέλεσης (σε πακέτα γραμμή TINY 5-10). Αποκτήστε καταπληκτικές προσαρμοστικές δυνατότητες με ένα SIMPLE σύνταξη!

Βρείτε την πηγή εδώ: https://github.com/scottgal/mostly clearning.atoms/blob/main/mostly clearning.ephemeral/sc/mostly clearning.ephemeral.complete

Ηγούμενοι

Διαβάστε το [προηγούμενο μέρος σχετικά με τα σήματα ]για κάποια διορατικότητα στις χρήσεις του. Παρουσιάζεται εδώ είναι το Readme.md από το πιο διαυγή.ephemeral.πλήρη pacakage που περιέχει τόσο τον πυρήνα ως επί το πλείστον διαυγή.ephemeral πακέτο και όλα τα πρότυπα, 'άτομα' (συντονιστές κ.λπ.) πακέτα σε ένα βολικό DLL.

Πυρήνας

Ή χρησιμοποιήστε τον πυρήνα ως επί το πλείστον διαυγής. ephemeral ένα TINY (κυριολεκτικά 10 τάξεις) που σας δίνει όλες τις πρώτες λειτουργίες.

Χαρακτηριστικά & DI

Ή αν θέλετε πλήρες χαρακτηριστικό βασισμένο async routing με απλή [EphemeralJob] και υπηρεσία.Ο συντονιστής σας χρησιμοποιεί την εγγραφή στυλ ως επί το πλείστον διαυγή. epemeral.attributes πακέτο.

Αυτό είναι πιθανώς ΤΟ θέμα του blog μου πηγαίνει προς τα εμπρός ... έχετε προειδοποιηθεί

Κυρίως διαυγής.Εφημερίς.Πλήρης

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

Όλα ως επί το πλείστον διαυγής.Ephemeral σε ένα μόνο DLL - περιορίζεται η εκτέλεση async με συντονισμό με βάση το σήμα.

dotnet add package mostlylucid.ephemeral.complete

Αυτό το πακέτο συγκεντρώνει όλο τον πυρήνα, το άτομο και τον κώδικα μοτίβο σε ένα σύνολο. Το τμήμα παρακάτω.


Πίνακας Περιεχομένων


Γρήγορη εκκίνηση

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

Το γνωστό services.AddCoordinator<T>() Βοηθοί και AddEphemeralSignalJobRunner<T>() Κρατήστε την εγγραφή υπηρεσιών συνοπτικό, αφήστε τον DI να κατέχει το νεροχύτη / runner, και κάντε το νέο ευθύνη / cache / logging ιστορίες ένα μόνο κλικ μακριά.

Δουλειές με γνώμονα τους στόχους

mostlylucid.ephemeral.complete Δέσμες mostlylucid.ephemeral.attributes, έτσι οι αγωγοί χαρακτηριστικών είναι μέρος του πυρήνα επιφάνεια. Αντιμετωπίστε τον δρομέα ως πρώτης θέσης καταναλωτή σήματος: διακοσμημένες μέθοδοι ενώνουν την ίδια caching, loging, και pinning ιστορίες, και κάθε χαρακτηριστικό μπορεί να δηλώσει Priority, επίπεδο εργασίας MaxConcurrency, Lane, Key πηγές, σήμα Εκπομπές, υπερβάσεις καρφιτσών/εκρήξεων, και επαναλήψεις.

Πληκτρολογημένα πλήκτρα γνωρίσματος:

  • Παραγγελίες & λωρίδες: Χρήση Priority, MaxConcurrency, και Lane να διατηρήσει την εργασία σε καθορισμένη τάξη ενώ hot μονοπάτια Μείνε χωριστά.
  • & Εντοπισμός πλήκτρων: OperationKey, KeyFromSignal, KeyFromPayload, και [KeySource] να σας βοηθήσει να συνεργαστείτε με την ομάδα σημαντικά κλειδιά για την καταγραφή, δίκαιο προγραμματισμό, και διαγνωστικά.
  • Τρύπημα & επανορθώσεις: Pin, ExpireAfterMs, AwaitSignals, MaxRetries, και RetryDelayMs Αφήστε τους χειριστές να επεκτείνουν την ορατότητα, την εκτέλεση της πύλης μέχρι να φτάσουν οι εξαρτήσεις, και να επουλωθούν με επαναλήψεις, ενώ εκπέμπουν σήματα αποτυχίας.
  • Χορογραφία σήματος: Emit EmitOnStart, EmitOnComplete, και EmitOnFailure για να σηματοδοτήσει κατάντη στάδια, log Παρατηρητές ή άλλοι συντονιστές χωρίς χειροκίνητη καλωδίωση.
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;
    }
}

Αυτός ο δρομέας κάθεται τώρα στην εκκίνηση και αντιδρά όποτε log.error.* ή οποιοδήποτε εκπέμπον σήμα χτυπά τον νεροχύτη. Οι χειριστές μπορούν επίσης να διαβάσουν κλειδιά από σήματα/πληρωμή, να δουλέψουν καρφίτσες μέχρι κατάντη acks, να εκπέμπουν σήματα ολοκλήρωσης/αποτυχίας, και υποδοχή σε λωρίδες για παραγγελία. Για DI-πρώτη χρήση ρυθμίσεων services.AddEphemeralSignalJobRunner<T>() (ή το πεδίο εφαρμογής) παραλλαγή) έτσι ο δρομέας και ο νεροχύτης διαχειρίζονται από το δοχείο.

[EphemeralJobs(SignalPrefix = "stage," DefaultLane = "pipeline")] Δημόσια σφραγισμένη τάξη StageJobs { [EphemeraralJob("ingest," EmitOn Complete = νέα[{ "stage.ingest.done" }] δημόσια εργασία IngestAsync(SignalEvent vt) => Console.Out.WriteLineAsync(evt.Signal) ·

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

}

var stageSink = new SignalSink() · περιμένουν χρησιμοποιώντας var στάδιοRunner = νέα EphemeralsignalJobRunner(stageSink, νέα[] {νέο στάδιο jobs() }) · στάδιοSink.Raise("stage.ingest") ·

Οι βαριές δουλειές μπορούν να βασιστούν σε ResponsibilitySignalManager.PinUntilQueried (προκαθορισμένο σχήμα ack responsibility.ack.*) έως να διατηρούν τις λειτουργίες τους ορατές μέχρις ότου ο μεταγενέστερος αναγνώστης προσγειώσει το ωφέλιμο φορτίο, ενώ OperationEchoMaker/ OperationEchoAtom επιμένουν στην τελική ροή σήματος έτσι ώστε οι ελεγκτές ή τα μόρια να μπορούν ακόμα να δοκιμάσουν την τελευταία κατάσταση ακόμη και μετά Το άτομο πεθαίνει.

Προγραμματισμένα καθήκοντα

mostlylucid.ephemeral.complete περιέχει επίσης mostlylucid.ephemeral.atoms.scheduledtasksΟρίστε το κορώνα ή το JSON προγράμματα μέσω ScheduledTaskDefinition (Cron, σήμα, προαιρετικό key, payload, description, timeZone, format, runOnStartup, κ.λπ.), και ScheduledTasksAtom enqueue ανθεκτικές εργασίες μέσω DurableTaskAtomΚάθε προγραμματισμένη δουλειά. σηκώνει το ρυθμισμένο σήμα μέσα σε ένα παράθυρο συντονιστή, έτσι ώστε να κληρονομήσει καρφίτσα, καταγραφή, και την ευθύνη σημασιολογία ενώ τα μόριά σας ή οι αγωγοί χαρακτηριστικών ανταποκρίνονται στο εκπέμπον κύμα σήματος.

Κάθε DurableTask μεταφέρει το χρονοδιάγραμμα Name, Signal, προαιρετική Key, ακόμη και ένα δακτυλογραφημένο Payload, και Description, έτσι ώστε κατάντη ακροατές γνωρίζουν αμέσως ποια εργασία έτρεξε και ποια μεταδεδομένα (στοιχεία, URL, κ.λπ.) να καταναλώνουν. DurableTaskAtom.WaitForIdleAsync() όταν απλά θέλετε να περιμένετε την τρέχουσα έκρηξη της προγραμματισμένης εργασίας να τελειώσει χωρίς να ολοκληρώσετε το άτομο, κρατώντας τον προγραμματιστή έτοιμο για το επόμενο cron τσιμπούρι.

Καταγραφή & σημάτων

mostlylucid.ephemeral.logging καθρέφτες Microsoft.Extensions.Logging σε σήματα και το αντίστροφο. SignalLoggerProvider στο εργοστάσιο logger σας έτσι log events αυξήσει log.* σήματα και άγκιστρα SignalToLoggerAdapter εάν Θέλετε τα σήματα να ρέουν πίσω στον τυπικό αγωγό καταγραφής.

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

Χρήση SignalToLoggerAdapter να καθρεφτίσει τα σήματα που προκύπτουν πίσω στα τυποποιημένα αρχεία καταγραφής έτσι ώστε η στοίβα παρακολούθησης σας βλέπει και τα δύο πλαϊνές πλευρές της γέφυρας.


Κύριοι συντονιστές

Συσκευασία: ως επί το πλείστον διαυγής. ephemeral

Ephemeral WorkCoordinator<Τ>

Μακροζωής ουρά εργασίας με οριοθετημένο νόμισμα και παρατηρήσιμο παράθυρο.

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

EphemeralKeyedWorkCoordinator<TKey, T>

Ανά κλειδί διαδοχική επεξεργασία - αντικείμενα με το ίδιο κλειδί που έχουν υποστεί επεξεργασία με τη σειρά τους.

await using var coordinator = new EphemeralKeyedWorkCoordinator<Order, string>(
    order => order.CustomerId,  // Key selector
    async (order, ct) => await ProcessOrder(order, ct),
    new EphemeralOptions
    {
        MaxConcurrency = 16,      // Total parallel
        MaxConcurrencyPerKey = 1  // Sequential per customer
    });

await coordinator.EnqueueAsync(order);

EphemeralResultCoordinator<TInput, TResult>

Η σύλληψη προκύπτει από ασύγχρονες επιχειρήσεις.

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

Συντονιστής εργασίας προτεραιότητας<Τ>

Πολλαπλές λωρίδες προτεραιότητας με ρυθμιζόμενο νόμισμα ανά λωρίδα.

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

Διαμόρφωση (EphemeralOptions)

new EphemeralOptions
{
    // Concurrency
    MaxConcurrency = 8,                    // Max parallel operations
    MaxConcurrencyPerKey = 1,              // For keyed coordinators
    EnableDynamicConcurrency = false,      // Allow runtime adjustment

    // Memory
    MaxTrackedOperations = 200,            // Window size (LRU eviction)
    MaxOperationLifetime = TimeSpan.FromMinutes(5),

    // Fair scheduling (keyed only)
    EnableFairScheduling = false,          // Prevent hot key starvation
    FairSchedulingThreshold = 10,

    // Signals
    Signals = sharedSink,                  // Shared signal sink
    OnSignal = evt => { },                 // Sync callback
    OnSignalAsync = async (evt, ct) => { }, // Async callback
    CancelOnSignals = new HashSet<string> { "circuit-open" },
    DeferOnSignals = new HashSet<string> { "backpressure" },
    DeferCheckInterval = TimeSpan.FromMilliseconds(100),
    MaxDeferAttempts = 50,

    // Signal handler limits
    MaxConcurrentSignalHandlers = 4,
    MaxQueuedSignals = 1000
}

Σήματα

Οι λειτουργίες εκπέμπουν σήματα για την οριζόντια παρατηρησιμότητα.

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

Σήματα και οριστικοποίηση ευθυνών

Πρέπει να κρατάμε τα αποτελέσματα ορατά αρκετά για τους μεταγενέστερους καταναλωτές; ResponsibilitySignalManager Σας επιτρέπει να καρφώσετε ένα λειτουργία έως ότου φτάσει ένα σήμα Ack (προκαθορισμένο μοτίβο responsibility.ack.* με το κλειδί=operationId). Παρέχετε ένα προαιρετική description έτσι η επιχείρηση μπορεί να περιγράψει την ευθύνη της, και να θέσει maxPinDuration Με χάρη Αυτόνομος αν ο καταναλωτής δεν εμφανιστεί ποτέ.

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 παραμένει μικρό (idivity, κλειδί, σήμα, χρονοσφραγίδα), ώστε να μπορείτε να καταγράψετε οποιαδήποτε ελάχιστη κατάσταση που σας ενδιαφέρει περίπου πριν από τη συλλογή της επιχείρησης.

Ο συντονιστής διατηρεί επίσης μια σύντομη ηχώ των τελικών σημάτων (ενήμερη μέσω EnableOperationEcho) ότι μπορείτε να επιθεωρήστε με GetEchoes() όταν χρειάζεται να ξαναπαίξετε το κομμένο κύμα σήματος χωρίς να κρατήσετε την πλήρη λειτουργία τριγύρω.

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 και OperationEchoCapacity να σας αφήσω να ισορροπήσετε πόσες ηχώ κρατάτε και πόσο καιρό παραμένουν, ώστε να μπορείτε να ξαναπαίξετε τις τελευταίες λέξεις απλά αρκετά μεγάλο χρονικό διάστημα για να επιφάνει διαγνωστικά.

Ο διαχειριστής ξεκλειδώνει αυτόματα όταν το ack πυροβολεί, αλλά μπορείτε να καλέσετε CompleteResponsibility(operationId) να τερματίσει το Υpiηρεσίε piου εξακολουθούν να αυξάνονται OperationFinalized όταν τα κόβει το παράθυρο, έτσι εγγραφείτε εάν θέλετε να εκπέμψετε ένα τελικό σήμα, διαγνωστικά καταγραφής, ή να εκτελέσετε τον καθαρισμό των τελευταίων λέξεων.


Άτομα (χτίσιμο μπλοκ)

Σταθερό έργοAtom

**Συσκευασία: ** κυρίως διαυγής.επεμφανής.άτομα. σταθερή εργασία

Σταθερή πισίνα εργαζομένων με στατιστικά. Minimal API περιτύλιγμα γύρω από EphemeralWorkCoordinator.

using Mostlylucid.Ephemeral.Atoms.FixedWork;

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

await atom.EnqueueAsync(item);

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

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

// Graceful shutdown
await atom.DrainAsync();

KeyedSequentialAtom

**Συσκευασία: ** ως επί το πλείστον διαυγή. epemeral.atoms.keyedsquential

Ανά κλειδί διαδοχική επεξεργασία με προαιρετικό δίκαιο προγραμματισμό.

using Mostlylucid.Ephemeral.Atoms.KeyedSequential;

await using var atom = new KeyedSequentialAtom<Order, string>(
    keySelector: order => order.CustomerId,
    body: async (order, ct) => await ProcessOrder(order, ct),
    maxConcurrency: 16,
    perKeyConcurrency: 1,           // Sequential per key
    enableFairScheduling: true);    // Prevent hot key starvation

await atom.EnqueueAsync(order1);  // Customer A
await atom.EnqueueAsync(order2);  // Customer A - waits for order1
await atom.EnqueueAsync(order3);  // Customer B - parallel with A

var (pending, active, completed, failed) = atom.Stats();
await atom.DrainAsync();

SignalAwareAtom

**Συσκευασία: ** ως επί το πλείστον διαυγή. epemeral.atoms.signalaware

Παύση ή ακύρωση εισαγωγής με βάση τα σήματα περιβάλλοντος.

using Mostlylucid.Ephemeral.Atoms.SignalAware;

var sink = new SignalSink();

await using var atom = new SignalAwareAtom<WorkItem>(
    async (item, ct) => await ProcessAsync(item, ct),
    cancelOn: new HashSet<string> { "shutdown", "circuit-open" },
    deferOn: new HashSet<string> { "backpressure.*" },
    deferInterval: TimeSpan.FromMilliseconds(100),
    maxDeferAttempts: 50,
    signals: sink,
    maxConcurrency: 8);

// Enqueue work
await atom.EnqueueAsync(item);

// Raise ambient signals
atom.Raise("backpressure.downstream");  // New items defer
sink.Raise("shutdown");                  // New items rejected (returns -1)

await atom.DrainAsync();

BatchingAtom

**Συσκευασία: ** ως επί το πλείστον διαυγή. epemeral.atoms.battching

Συλλέξτε αντικείμενα σε παρτίδες κατά μέγεθος ή χρονικό διάστημα.

using Mostlylucid.Ephemeral.Atoms.Batching;

await using var atom = new BatchingAtom<LogEntry>(
    onBatch: async (batch, ct) =>
    {
        Console.WriteLine($"Flushing {batch.Count} entries");
        await FlushToDatabase(batch, ct);
    },
    maxBatchSize: 100,
    flushInterval: TimeSpan.FromSeconds(5));

// Items are batched automatically
atom.Enqueue(new LogEntry("User logged in"));
atom.Enqueue(new LogEntry("Request received"));
// ... batch flushes when full OR after 5 seconds

RetryAtom

Συσκευασία: ως επί το πλείστον διαυγή. epemeral.atoms.retry

Επεκτατικό περιτύλιγμα αντιγράφων ασφαλείας.

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

Άτομα αποθήκευσης δεδομένων

Συσκευασία: ως επί το πλείστον διαυγή. epemeral.atoms.data

Κοινόχρηστες ρυθμίσεις για άτομα αποθήκευσης (DataStorageConfig, IDataStorageAtom<TKey, TValue>) plus the signal conventions that drive file, SQLite, and PostgreSQL adapters.

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

Χρήση της ίδιας DataStorageConfig με Mostlylucid.Ephemeral.Atoms.Data.Sqlite ή Mostlylucid.Ephemeral.Atoms.Data.Postgres Εφαρμογές για διαρκή, με γνώμονα το σήμα επιμονή που τροφοδοτείται από SQLite/Postgres. saved.data.{dbname} σήματα για την εκκίνηση κατάντη εργασιών, ενώ load.data.{dbname} ενεργοποιεί ενυδατικές κρύπτες.


MoleculeRunner & AtomTrigger

**Συσκευασία: ** ως επί το πλείστον διαυγή. epemeral.atoms.molecules

Εκτυπώσεις με blueprints MoleculeBlueprintBuilder σας επιτρέπει να καθορίσετε τα άτομα (πληρωμή, απογραφή, αποστολή, κοινοποίηση) που θα πρέπει να εκτελείται όταν ένα σήμα, όπως order.placed Φτάνει. MoleculeRunner Ακούει για τη σκανδάλη. μοτίβο, δημιουργεί ένα κοινό MoleculeContextκαι εκτελεί κάθε βήμα ενώ εγγράφεστε για να ξεκινήσετε/ολοκληρώσετε τα γεγονότα. AtomTrigger όταν το σήμα ενός ατόμου θα πρέπει να ξεκινήσει έναν άλλο συντονιστή ή μόριο.

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

Τα βήματα του μορίου μπορούν να ανεβάσουν πρόσθετα σήματα (ctx.Raise("order.shipping.start")) έτσι το υπόλοιπο του συστήματος παίρνει το Μπάτον.


SlidingCacheAtom

**Συσκευασία: ** ως επί το πλείστον διαυγή. epemeral.atoms.slidingcache

Cache με συρόμενη λήξη - πρόσβαση σε ένα αποτέλεσμα επαναφέρει 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}");

EphemeralLruCache

Συσκευασία: Πυρήνας (mostlylucid.ephemeral) Αυτό-optimizing cache με συρόμενο TTL σε κάθε χτύπημα και επέκταση TTL για Καυτά κλειδιά.

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

Συμβουλή: MemoryCache μπορεί να ρυθμιστεί για συρόμενη λήξη, αλλά ποτέ δεν εκπέμπει τα ζεστά / ψυχρά σήματα ή επεκτείνει TTL για ζεστά κλειδιά. EphemeralLruCache είναι η αυτο-βελτιωτική προεπιλογή στο βασικό πακέτο (και σε SqliteSingleWriter) Όποτε θέλετε η κρύπτη να επικεντρωθεί στο ενεργό σύνολο εργασίας.

Echo Maker

Συσκευασία: ως επί το πλείστον διαυγής. ephemeral.atoms.eco

Συλλάβετε τις δακτυλογραφημένες λέξεις που εκπέμπονται πριν από το κούρεμα. ωφέλιμα φορτία (αντιστοιχούν σε ωφέλιμο φορτίο) ActivationSignalPattern / CaptureSignalPattern) και πότε OperationFinalized πυρκαϊές που παράγει OperationEchoEntry<TPayload> αρχεία που μπορείτε να επιμείνετε μέσω 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");

Οι θέσεις εργασίας αποτελούν απλώς το δακτυλογραφημένο σήμα με οποιαδήποτε κατάσταση θεωρούν κρίσιμη, και ο κατασκευαστής διατηρεί το σύνολο εργασίας Περιορίζεσαι ενώ επιμένεις στην ηχώ.


Μοτίβα (Ready-to-Use)

SignalBasedCircuitBreaker

**Συσκευασία: ** ως επί το πλείστον διαυγής. epemeral.patterns.curcreambreaker

Ανεπίσημος διακόπτης κυκλώματος χρησιμοποιώντας παράθυρο ιστορικού σήματος.

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

SignalDrivenBackpressure

**Συσκευασία: ** ως επί το πλείστον διαυγή. epemeral.patterns. backpressure

Διαχείριση βάθους αναμονής με αυτόματη αναβολή στα σήματα πίσω πίεσης.

using Mostlylucid.Ephemeral.Patterns.Backpressure;

var sink = new SignalSink();

await using var coordinator = SignalDrivenBackpressure.Create<WorkItem>(
    async (item, ct) => await ProcessAsync(item, ct),
    sink,
    maxConcurrency: 4);

// Enqueue work
await coordinator.EnqueueAsync(item);

// When downstream is slow
sink.Raise("backpressure.downstream");  // New work auto-defers

// When recovered
sink.Retract("backpressure.downstream"); // Work resumes

ControlledFanOut

**Συσκευασία: ** ως επί το πλείστον διαυγή. epemeral.patterns.controlledfanout

Global + per-key gating για ελεγχόμενο παραλληλισμό.

using Mostlylucid.Ephemeral.Patterns.ControlledFanOut;

await using var fanout = new ControlledFanOut<string, Request>(
    keySelector: req => req.TenantId,
    body: async (req, ct) => await ProcessAsync(req, ct),
    maxGlobalConcurrency: 100,  // Total parallel across all tenants
    perKeyConcurrency: 5);      // Max 5 parallel per tenant

// Items for same tenant processed with limit
await fanout.EnqueueAsync(requestA);  // Tenant1
await fanout.EnqueueAsync(requestB);  // Tenant1 - waits if 5 already running
await fanout.EnqueueAsync(requestC);  // Tenant2 - parallel with Tenant1

await fanout.DrainAsync();

AdaptiveRateService

**Συσκευασία: ** ως επί το πλείστον διαυγής. epemeral. patterns. adaptiverate

Ρυθμός μετάδοσης σήματος που περιορίζεται με αυτόματη οπισθοδρόμηση.

using Mostlylucid.Ephemeral.Patterns.AdaptiveRate;

await using var service = new AdaptiveRateService<ApiRequest>(
    async (req, ct) => await CallApiAsync(req, ct),
    maxConcurrency: 8);

// Process with automatic rate limit handling
await service.ProcessAsync(request);

// When API returns 429, emit signal with retry-after
// Signal: "rate-limit:500ms"
// Service auto-parses and delays

Console.WriteLine($"Pending: {service.PendingCount}, Active: {service.ActiveCount}");

DynamicConcurrencyDemo

**Συσκευασία: ** ως επί το πλείστον διαυγή. epemeral.patterns. Dynamicconcurrency

Η κλιμάκωση του χρόνου λειτουργίας βασίζεται σε σήματα φορτίου.

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

Keyed PriorityFanOut

**Συσκευασία: ** ως επί το πλείστον διαυγής. epemeral.patterns.keyed primorityfanout

Προτεραιότητα λωρίδες με ανά κλειδί παραγγελία διατηρείται.

using Mostlylucid.Ephemeral.Patterns.KeyedPriorityFanOut;

await using var fanout = new KeyedPriorityFanOut<string, UserCommand>(
    keySelector: cmd => cmd.UserId,
    body: async (cmd, ct) => await HandleCommand(cmd, ct),
    maxConcurrency: 32,
    perKeyConcurrency: 1,  // Sequential per user
    maxPriorityDepth: 100);

// Normal lane
await fanout.EnqueueAsync(normalCommand);

// Priority lane - jumps the queue for that user
bool accepted = await fanout.EnqueuePriorityAsync(urgentCommand);

// Check lane depths
var counts = fanout.PendingCounts;
Console.WriteLine($"Priority: {counts.Priority}, Normal: {counts.Normal}");

await fanout.DrainAsync();

ReactiveFanOutPipeline

**Συσκευασία: ** ως επί το πλείστον διαυγή. epemeral. patterns. activatefanout

Αγωγός δύο σταδίων με αυτόματη πίεση πλάτης.

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

Ανιχνευτής σημάτωνAnomalyName

**Συσκευασία: ** ως επί το πλείστον διαυγής. epemeral.patterns.anomalydetector

Ανίχνευση ανωμαλίας παραθύρου.

using Mostlylucid.Ephemeral.Patterns.AnomalyDetector;

var sink = new SignalSink();

var detector = new SignalAnomalyDetector(
    sink,
    pattern: "error.*",
    threshold: 5,
    window: TimeSpan.FromSeconds(10));

// Check for anomalies
if (detector.IsAnomalous())
{
    Console.WriteLine("Anomaly detected! Too many errors.");
    TriggerAlert();
}

// Get current match count
int errorCount = detector.GetMatchCount();
Console.WriteLine($"Errors in window: {errorCount}");

SignalCoordinatedReads

**Συσκευασία: ** ως επί το πλείστον διαυγή. epemeral.patterns.signal coordinaryreads

Ησυχία διαβάζει κατά τη διάρκεια των ενημερώσεων χωρίς σκληρές κλειδαριές.

using Mostlylucid.Ephemeral.Patterns.SignalCoordinatedReads;

// Run demo: readers pause when update signal is present
var result = await SignalCoordinatedReads.RunAsync(
    readCount: 10,
    updateCount: 1);

Console.WriteLine($"Reads: {result.ReadsCompleted}, Updates: {result.UpdatesCompleted}");
Console.WriteLine($"Signals: {string.Join(", ", result.Signals)}");

// Manual implementation:
var sink = new SignalSink();

await using var readers = new EphemeralWorkCoordinator<Query>(
    body,
    new EphemeralOptions
    {
        DeferOnSignals = new HashSet<string> { "update.in-progress" },
        Signals = sink
    });

// Readers auto-defer when update is running
sink.Raise("update.in-progress");  // Readers wait
sink.Raise("update.done");         // Readers resume

SignalingHttpClient

**Συσκευασία: ** ως επί το πλείστον διαυγής.ephemeral.patterns.signalinghttp

Πελάτης HTTP με σήματα προόδου.

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

SignalLogWatcherName

**Συσκευασία: ** ως επί το πλείστον διαυγής.ephemeral.patterns.signallogwatcher

Παρακολουθήστε παράθυρο σήματος για μοτίβα και ενεργοποιήστε 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)

ΤηλεμετρίαSignalHandler

**Συσκευασία: ** ως επί το πλείστον διαυγής. epemeral.patterns.telemetry

Η ενσωμάτωση OpenTelemetry/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();

LongWindowDemo

**Συσκευασία: ** ως επί το πλείστον διαυγή. epemeral.patterns. longwindowdemo

Διαδηλώνει μεγάλες ρυθμίσεις παραθύρων για διαδρομές ελέγχου.

using Mostlylucid.Ephemeral.Patterns.LongWindowDemo;

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

SignalReactionShowcase

**Συσκευασία: ** Κυρίως διαυγής. epemeral.patterns.signalreactionsshowcase

Διαμαρτύρεται για τα μοτίβα αποστολής σήματος και τις επανακλήσεις.

using Mostlylucid.Ephemeral.Patterns.SignalReactionShowcase;

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

ΣυνεχώςSignalWindowName

**Συσκευασία: ** ως επί το πλείστον διαυγή. epemeral.patterns. continent window

Παράθυρο σήματος με SQLite επιμονή - επιβιώνει διαδικασία επανεκκινεί.

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

Οι σύγχρονες ρίζες DI μπορεί να προτιμούν τους μικρότερους βοηθούς, όπως services.AddCoordinator<T>(...), services.AddScopedCoordinator<T>(...), ή services.AddKeyedCoordinator<T, TKey>(...) Από τότε που διαβάζουν κανονικά. AddX Εγγραφές · απλώς αναθέτουν στους ειδικούς βοηθούς της Εφέδρας κάτω από την κουκούλα.


Πλαίσια-στόχοι

  • NET 6.0, 7.0, 8.0, 9.0, 10.0

Άδεια

Μη απαλλαγή (δημόσιος τομέας)

Finding related posts...
logo

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