Ένα μικροσκοπικό πρωτόγονο που μετατρέπει την ταυτόχρονη εργασία σε ένα συντονισμένο, προσαρμοστικό σύστημα.
"Το μοτίβο Εφέμεραλ Σήματα"
Το Μέρος 1 Κατασκευάσαμε εφήμερη εκτέλεση - περιορισμένη, ιδιωτική, αυτοκαθαριζόμενη ασύγχρονη ροή εργασίας. Μέρος 2 Το μετατρέψαμε σε μια επαναχρησιμοποίητη βιβλιοθήκη με συντονιστές, αγωγούς κλειδιού και ενσωμάτωση DI.
Αυτό το άρθρο προσθέτει ένα μικρό χαρακτηριστικό που αλλάζει τα πάντα: σήματα.
Αυτό είναι τώρα στο πιο διαυγή.ephemerals πακέτο Nuget επίσης Περισσότερα από 20 διαυγής.εφημέρια μοτίβα και "άτομα".
Ο πλήρης πηγαίος κώδικας είναι στο ως επί το πλείστον διαυγή.άτομα Αποθετήριο GitHub
Η υποδομή του σήματος ζει σε:
Αρχείο > Σκοπός > > > > Αρχείο > > > > > > > > > > > > > > > > > > > > > > < > > < > > > < > < > > < > > > < > > > > > < > > > > > > < < > > > > > > > > < > > < > > < < > > > > < > > > > > > < < < < > < < < < < < < < < < <
| ------ | --------- |
|---|---|
| Ephemeraral Operation.cs ~ Εκπομπή σήματος και ανάκληση από τις λειτουργίες ~ | |
EphemeralOptions.cs Διαμόρφωση αντιδραστικού σήματος (CancelOnSignals, DeferOnSignals, OnSignal, OnSignalRetracted) |
|
| StringPatternMatcher.cs Το μοτίβο Glob-style ταιριάζει για το φιλτράρισμα σημάτων | |
SignalDispatcher.cs Ροή σήματος Async με ταίριασμα μοτίβο (υποστηρίγματα) *, ?, λίστες κόμμα, καθορισμένη τάξη) |
|
| Παραδείγματα/SignalingHttpClient.cs Το δείγμα εκπομπών σήματος με λεπτό γρανάζιο για τις κλήσεις HTTP | |
| Παραδείγματα/Προσαρμοσμένη Υπηρεσία Μετάφρασης. cs Προσαρμοζόμενος ρυθμός περιορισμού με την αναβολή με βάση το σήμα | |
| Παραδείγματα/SignalBasedCircuitBreaker.cs Ο διακόπτης κυκλώματος διαβάζει εφήμερο παράθυρο σήματος | |
| Παραδείγματα/ΤηλεμετρίαSignalHandler.cs Επεξεργασία σήματος Async με ενσωμάτωση τηλεμετρίας |
Οι εφήμεροι συντονιστές μας είναι καλοί στην επεξεργασία των εργασιών, αλλά είναι απομονωμένοι.Κάθε συντονιστής γνωρίζει για τις δικές του λειτουργίες, αλλά δεν έχει καμία επίγνωση για το τι συμβαίνει αλλού στο σύστημα.
// Translation coordinator has no idea that...
await translationCoordinator.EnqueueAsync(request);
// ...the API just hit a rate limit
// ...another service is experiencing backpressure
// ...a downstream dependency is slow
Θα μπορούσαμε να καλωδιώσουμε σαφείς εξαρτήσεις, αλλά αυτό δημιουργεί σύζευξη. ευαισθητοποίηση του περιβάλλοντος - συντονιστές που μπορούν να αισθανθούν το περιβάλλον τους χωρίς να είναι άμεσα συνδεδεμένοι.
Τα σήματα αφήνουν τα άτομα της εκτέλεσης να αφήνουν ίχνη στο εφήμερο παράθυρό τους.
Οι συντονιστές μπορούν στη συνέχεια να αλλάξουν τη συμπεριφορά τους με βάση τα σήματα ορατά στο παράθυρό τους.
Από Σήματα. cs:
public readonly record struct SignalEvent(
string Signal,
long OperationId,
string? Key,
DateTimeOffset Timestamp,
SignalPropagation? Propagation = null)
{
public int Depth => Propagation?.Depth ?? 0;
public bool WouldCycle(string signal) => Propagation?.Contains(signal) == true;
public bool Is(string name) => Signal == name;
public bool StartsWith(string prefix) => Signal.StartsWith(prefix, StringComparison.Ordinal);
}
Τα σήματα αυτά ζουν στο εφήμερο παράθυρο δίπλα στην εγχείρηση.
Δεν υπάρχει μεσίτης μηνυμάτων, δεν υπάρχει ξεχωριστή υποδομή, απλά δεσμεύσεις που συνδέονται με τις επιχειρήσεις.
Επειδή τα σήματα ζουν μέσα στο εφήμερο παράθυρο, κληρονομούν τις εγγυήσεις του: περιορισμένο μέγεθος, αυτόματη γήρανση και μηδενικό κύκλο ζωής από πάνω.
Ένα σήμα, από μόνο του, δεν προκαλεί εκτέλεση. Καταγράφει μόνο ένα γεγονός στο εφήμερο παράθυρο. Τίποτα δεν τρέχει επειδή εκπέμπει ένα σήμα.
Εάν ένας συντονιστής ορίζει έναν χειριστή OnSignal, λειτουργεί συγχρονικά όταν εκπέμπεται ένα σήμα - αλλά μόνο επειδή ο συντονιστής επέλεξε να το επισυνάψει. Οι Έμιτερς δεν ξέρουν ή δεν νοιάζονται. Η αφαίρεση όλων των χειριστών αφήνει αμετάβλητη τη συμπεριφορά του πυρήνα.
Ένα σήμα συνδέεται μόνο με το άτομο/τη λειτουργία που το εκπέμπει. Κανένα σήμα δεν μεταλλάσσεται ή δεν σχολιάζει άλλο άτομο. Δεν υπάρχει κοινόχρηστο λεωφορείο.
Η εγκατάλειψη ενός σήματος προσθέτει ένα γεγονός στην ιστορία του ατόμου. Τα σήματα δεν ενημερώνονται ή δεν αντιγράφονται ποτέ. Αφαιρούν μόνο τα σήματα του πομπού.
Τα σήματα υπάρχουν μόνο μέσα στο παράθυρο εφήμερο του συντονιστή. Λήγουν αυτόματα καθώς το παράθυρο γερνάει. Τίποτα δεν παραμένει εκτός αν χτίσεις ρητά επιμονή.
Όταν ένας συντονιστής ελέγχει τα σήματα για το κλειδί ΚΚ, σαρώνει:
τα άτομα στο παράθυρό του
και τα σήματα τοπικά σε αυτά τα άτομα Οι παρατηρητές ποτέ δεν τροποποιούν μια κατάσταση σήματος ατόμων.
Εάν θέλετε σήματα για την οδήγηση async ροή εργασίας, θα πρέπει να χρησιμοποιήσετε:
SignalDispatcher
AsyncSignal Processor
ή άλλους προσαρμοστές.
Αυτά είναι προαιρετικά στρώματα πάνω από τα σήματα, όχι μέρος της σημασιολογίας τους.
Κανένα άτομο δεν μπορεί να αλλάξει το στιγμιότυπο, την κατάσταση, τα σήματα ή τα μεταδεδομένα ενός άλλου ατόμου. Ο συντονισμός πραγματοποιείται μέσω:
σήματα
ανίχνευση
παράθυρα
Πολιτικές
Όχι μέσω των γραπτών.
Ένα SignalSink ή συρματόπλεγμα επιφάνειες μπορεί να εμφανίσει μια συνδυασμένη άποψη των σημάτων Αλλά πάντα διαβάζεται μόνο, ποτέ δεν είναι έγκυρο, ποτέ δεν είναι γραμμένο.
Εάν αποσπαστούν όλοι οι χειριστές (OnSignal, αποστολές, επεξεργαστές), το σύστημα παραμένει πλήρως ορθό και προβλέψιμο. Τα σήματα έχουν ακόμα νόημα επειδή είναι γεγονότα, όχι πυροδοτήσεις.
Γεγονότα είναι σαν ένα τηλεφώνημα:
"Σε παίρνω τώρα, σήκωσέ το και αντέδρασε."
Σήματα είναι σαν πατημασιές στο χιόνι:
"Άφησα πατημασιές, αν θες να μάθεις που πήγα, κοίτα, αν δεν σε νοιάζει, αγνόησέ το."
Αυτός είναι ο λόγος για τον οποίο τα σήματα δεν σπάνε ποτέ, ποτέ δεν μπλοκάρουν, και ποτέ δεν αλληλεπιδρούν με τη ροή ελέγχου εκτός αν εσείς επιλογή να τους κάνω δημοσκοπήσεις.
Η διάκριση έχει σημασία επειδή τα σήματα εξαλείφουν:
□ Πρόβλημα Γεγονότα Έχουν Σημάδια Αποφεύγετε το |---------|:--------------:|:----------------:| Συνδρομητής → Κανένας Ο συγχρονισμός εξαρτήσεων η άμεση αντίδραση που απαιτείται περιβάλλον, δημοσκόπηση όταν είναι έτοιμη Οι χειριστές μπορούν να ενεργοποιήσουν τους χειριστές. Δεν υπάρχουν βρόχοι εκτός αν ζητηθεί ρητά. Εντολές σημασιολογίας θέματα παραγγελίας άσχετα Οι εγγυήσεις παράδοσης πρέπει να παραδοθούν/επεξεργάζονται όχι παράδοση, μόνο ύπαρξη Σφάλμα διάδοσης Σφάλμα διάδοσης Σφάλματα χειριστηρίου πολλαπλασιάζεται Απομονωμένος Κίνδυνοι αναστάτωσης Αστοχίες Κασκάντ Ένα χειριστής αποτυγχάνει, αλυσίδα σπάει Δεν αλυσίδα για να σπάσει
Χαρακτηριστικά γνωρίσματα |---------|--------|---------| Συνδυάζοντας ~ Strong (εκδότης → συνδρομητές) ~ Κανένας ~ Χρονοδιακόπτης - Άμεσο περιβάλλον Παράδοση Εγγυημένη/προστεθειμένη Δεν παράδοση, μόνο ύπαρξη Η αντίδραση είναι απαραίτητη Προαιρετική □ Πληρωμή δεδομένων □ Συχνά βαριά μεταδεδομένα μικροσκοπικών συμβολοσειρών Δεν υπάρχει παράθυρο LRU Slivering Ο χρόνος ζωής είναι στιγμιαίος Αυτόματη αποσυναρμολόγηση Σχεδόν καθόλου
Προσέγγιση γεγονότων (κλασική):
public event Action RateLimited;
try
{
await CallApiAsync();
}
catch (RateLimitException)
{
RateLimited?.Invoke(); // Makes someone else act right now
}
Προβλήματα:
Προσέγγιση σήματος (επέμφαση):
try
{
await CallApiAsync(ct);
}
catch (RateLimitException ex)
{
op.Signal($"rate-limit:{ex.RetryAfterMs}ms");
throw;
}
Αργότερα, κάπου εντελώς διαφορετικά.
if (translator.HasSignal("rate-limit"))
{
await Task.Delay(1000); // Act when *we* choose
}
Σημειώσεις:
Έλεγχος μεταφοράς γεγονότων, πλαίσιο μεταφοράς σημάτων.
Αυτό είναι όλο το νοητικό μοντέλο σε μία πρόταση.
Ανάλογα με το παρελθόν σας:
Ορίστε το κοινό |----------|------------| | Τα Συστήματα Σκέφτονται Ένα στιγμινεργικό υπόστρωμα για έμμεσο συντονισμό | Μηχανικοί Τα μεταδεδομένα ελαφρού βάρους που συνδέονται με τις λειτουργίες σε ένα ρυθμιζόμενο συρόμενο παράθυρο | PL/Σπασίκλας νομίσματος Ένας έμμεσος, χρονικός μαυροπίνακας που συνδυάζεται με οριοθετημένη σημασιολογία concurrency | Χρήστες πλαισίου Μια αυτοκαθαριζόμενη επιφάνεια που μπορείτε να ρωτήσετε ανά πάσα στιγμή
Τα σήματα είναι ίχνη που έχουν απομείνει σε μια κοινή, οριοθετημένη επιφάνεια μνήμης. στιγμινεργική μοντέλο συντονισμού - το ίδιο ένα μυρμήγκια χρήση, το ίδιο ένα σύστημα blackboard που χρησιμοποιείται στην αρχή της AI, και το ίδιο ένα σύγχρονο δίκτυο κουτσομπολιού CRDT υπαινίσσεται.
using var activity = source.StartActivity("ProcessOrder");
activity?.SetTag("order.id", orderId);
activity?.SetTag("rate.limited", true);
Το καλύτερο για: Διανεμημένος εντοπισμός σε όλες τις υπηρεσίες, μακροχρόνια αποθήκευση τηλεμετρίας, ταυτότητες συσχέτισης.
Χρήση τηλεμετρίας όταν: Πρέπει να ανιχνεύσετε αιτήματα σε πολλαπλές υπηρεσίες, να αποθηκεύσετε μετρήσεις για ανάλυση, ή να ενσωματώσετε με εργαλεία παρακολούθησης.
Χρησιμοποιήστε τα σήματα Ephemerals όταν: Χρειάζεστε κατά τη διαδικασία ευαισθητοποίηση του περιβάλλοντος, αντιδραστικό συντονισμό, ή δεν θέλετε την υποδομή τηλεμετρίας.
var rateLimits = Observable.FromEventPattern<RateLimitEventArgs>(
h => api.RateLimitHit += h,
h => api.RateLimitHit -= h);
rateLimits
.Throttle(TimeSpan.FromSeconds(1))
.Subscribe(e => HandleRateLimit(e));
Το καλύτερο για: Σύνθετη επεξεργασία γεγονότων, χρονοβόρες εργασίες, συνδυάζοντας πολλαπλές ροές γεγονότων.
Χρήση Rx όταν: Χρειάζεστε περίπλοκες χρονικές ερωτήσεις (παράθυρο, αποκήρυξη, συνδυάζοντας ρεύματα).
Χρησιμοποιήστε τα σήματα Ephemerals όταν: Θέλετε απλούστερη ανίχνευση δημοσκοπήσεων, αυτόματο καθάρισμα, ή ενσωμάτωση με παρακολούθηση λειτουργίας.
public class RateLimitNotification : INotification
{
public int RetryAfterMs { get; init; }
}
await _mediator.Publish(new RateLimitNotification { RetryAfterMs = 5000 });
Το καλύτερο για: Αποσυνδεδεμένος χειρισμός συμβάντων κατά τη διαδικασία με πολλούς χειριστές.
Χρήση MediatR όταν: Θέλετε πολλούς χειριστές να αντιδρούν στο ίδιο γεγονός συγχρονικά.
Χρησιμοποιήστε τα σήματα Ephemerals όταν: Θέλετε την αντίληψη του περιβάλλοντος χωρίς ρητή συνδρομή, ιστορικό αυτοκαθαρισμού, ή ενσωμάτωση με περιορισμένη εκτέλεση.
var circuitBreaker = Policy
.Handle<HttpRequestException>()
.CircuitBreakerAsync(5, TimeSpan.FromSeconds(30));
Το καλύτερο για: Διαμονή γύρω από ατομικές κλήσεις με αυτόματη διαχείριση του κράτους.
Χρησιμοποιήστε την Polly όταν: Χρειάζεστε αντοχή ανά κλήση με αυτόματη ημιανοικτή/κλειστή μετάβαση.
Χρησιμοποιήστε τα σήματα Ephemerals όταν: Θέλετε την ευαισθητοποίηση του περιβάλλοντος σε πολλές λειτουργίες, συνήθεια λογική κύκλωμα, ή ενσωμάτωση με την παρακολούθηση λειτουργίας.
Συνδυάστε τα.: Χρησιμοποιήστε Polly μέσα στο σώμα εργασίας σας, εκπέμπουν σήματα όταν τα κυκλώματα εκτροχιάζονται.
Προσεγγιση αυτοκαθαρισμου Ανίχνευση περιβάλλοντος Αποσυνδεδεμενη Λογική συνήθειας Ενσωμάτωση |----------|:-------------:|:---------------:|:---------:|:------------:|:-----------:| Ανοιχτή Τηλεμετρία ~ Reactive Extensions - Συγκρότημα~ Μεσίτης Πόλι Circuit Breaker . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . | Εφήμερα Σήματα Ας το κάνουμε αυτό.
Εκτέλεση πράξεων ISignalEmitter:
public interface ISignalEmitter
{
// Emit signals
void Emit(string signal);
bool EmitCaused(string signal, SignalPropagation? cause);
// Retract (remove) signals
bool Retract(string signal);
int RetractMatching(string pattern);
bool HasSignal(string signal);
long OperationId { get; }
string? Key { get; }
}
Μέσα στο σώμα της εργασίας σας:
await coordinator.ProcessAsync(async (item, op, ct) =>
{
try
{
var result = await CallExternalApiAsync(item, ct);
if (result.WasCached)
op.Signal("cache-hit");
if (result.Duration > TimeSpan.FromSeconds(2))
op.Signal("slow-response");
}
catch (RateLimitException ex)
{
op.Signal("rate-limit");
op.Signal($"rate-limit:{ex.RetryAfterMs}ms");
throw;
}
catch (TimeoutException)
{
op.Signal("timeout");
throw;
}
});
Τα σήματα είναι απλά χορδές. Χρησιμοποιήστε απλά ονόματα ("rate-limit") ή δομημένες ονομασίες ("rate-limit:5000ms").
Τα φίλτρα μοτίβο χρησιμοποιούν σημασιολογία glob (*, ?) και λίστες comma υποστήριξης ("error.*,timeout"). Ταίριασμα είναι αποφασιστικής σημασίας και κατανομή-φως μέσω StringPatternMatcher.
Για πολύ λεπτομερή παρατηρητικότητα, μπορείτε να εκπέμψετε σήματα σε κάθε στάδιο μιας επιχείρησης. SignalingHttpClient που αποδεικνύει αυτό το μοτίβο:
using Mostlylucid.Helpers.Ephemeral.Examples;
// Inside your work body where you have access to the operation's emitter:
await coordinator.ProcessAsync(async (request, op, ct) =>
{
var data = await SignalingHttpClient.DownloadWithSignalsAsync(
httpClient,
new HttpRequestMessage(HttpMethod.Get, request.Url),
op, // ISignalEmitter
ct);
// Process the downloaded data...
});
Αυτό εκπέμπει σήματα σε κάθε στάδιο:
Σήμα όταν
|--------|------|
| stage.starting Πριν από την έναρξη της αίτησης
| progress:0 Αρχικό σηάδι piροόδου
| stage.request Το αίτημα του HTTP που στάλθηκε
| stage.headers Η απάντηση στα κεφαλάρια που λάβαμε
| stage.reading Άρχισε να διαβάζει το σώμα
| progress:XX Ποσοστό προόδου (0- 100) κατά τη λήψη
| stage.completed Ολοκληρώθηκε η λήψη
Στη συνέχεια μπορείτε να ρωτήσετε αυτά με μοτίβο που ταιριάζει:
// Find all stage transitions
var stages = coordinator.GetSignalsByPattern("stage.*");
// Check download progress
var progress = coordinator.GetSignalsByPattern("progress:*");
// Check if any download is still in progress
if (coordinator.HasSignalMatching("stage.reading") &&
!coordinator.HasSignalMatching("stage.completed"))
{
// Download in progress
}
Αυτό είναι χρήσιμο για προσωρινά κράτη:
await coordinator.ProcessAsync(async (item, op, ct) =>
{
// Mark as processing
op.Emit("processing");
try
{
await ProcessItemAsync(item, ct);
// Success - retract the processing signal
op.Retract("processing");
op.Emit("completed");
}
catch (RetryableException)
{
// Keep processing signal, add retry info
op.Emit("retrying");
}
catch (Exception)
{
// Remove all temporary signals
op.RetractMatching("processing*");
op.Emit("failed");
throw;
}
});
Όπως ακριβώς και η εκπομπή σήματος, οι ανακλήσεις μπορούν να ενεργοποιήσουν τις επανακλήσεις:
var coordinator = new EphemeralWorkCoordinator<Request>(
body,
new EphemeralOptions
{
// Sync retraction handler
OnSignalRetracted = evt =>
{
_metrics.DecrementGauge(evt.Signal);
Console.WriteLine($"Signal {evt.Signal} retracted from op {evt.OperationId}");
if (evt.WasPatternMatch)
Console.WriteLine($" (matched pattern: {evt.Pattern})");
},
// Async retraction handler
OnSignalRetractedAsync = async (evt, ct) =>
{
await _telemetry.TrackRetraction(evt.Signal, evt.OperationId, ct);
}
});
Η SignalRetractedEvent περιλαμβάνει:
Signal - Το ανασυρμένο όνομα σήματοςOperationId - Η επιχείρηση που το απέσυρε.Key - Το κλειδί της επιχείρησης (εάν υπάρχει)Timestamp - Όταν συνέβη η ανάκλησηWasPatternMatch - Είναι αλήθεια αν ανασυρθεί μέσω RetractMatchingPattern - Το μοτίβο που χρησιμοποιείται (αν ταιριάζει μοτίβο)await coordinator.ProcessAsync(async (request, op, ct) =>
{
// Check if we already have a rate limit signal
if (op.HasSignal("rate-limited"))
{
// We're in recovery mode
await Task.Delay(1000, ct);
}
try
{
var response = await _api.SendAsync(request, ct);
// Success! Remove any rate limit signal
if (op.Retract("rate-limited"))
{
op.Emit("rate-limit-cleared");
}
}
catch (RateLimitException ex)
{
op.Emit("rate-limited");
op.Emit($"rate-limit:{ex.RetryAfterMs}ms");
throw;
}
});
Όλοι οι συντονιστές παρέχουν βελτιστοποιημένη αναζήτηση σήματος:
// Check if any recent operation hit a rate limit
if (coordinator.HasSignal("rate-limit"))
{
await Task.Delay(1000);
}
// Count slow responses in the window
var slowCount = coordinator.CountSignals("slow-response");
if (slowCount > 10)
{
await ThrottleAsync();
}
// Get signals by pattern
var httpErrors = coordinator.GetSignalsByPattern("http.error.*");
// Get signals since a time
var recentSignals = coordinator.GetSignalsSince(DateTimeOffset.UtcNow.AddMinutes(-1));
// Get signals for a specific key
var userSignals = coordinator.GetSignalsByKey("user-123");
Από EphemeralOptions.cs:
Οι συντονιστές μπορούν να αντιδρούν αυτόματα στα σήματα:
var coordinator = new EphemeralWorkCoordinator<Request>(
body,
new EphemeralOptions
{
// Cancel new work if these signals are present
CancelOnSignals = new HashSet<string> { "system-overload", "circuit-open" },
// Defer new work while these signals are present
DeferOnSignals = new HashSet<string> { "rate-limit" },
MaxDeferAttempts = 10,
DeferCheckInterval = TimeSpan.FromMilliseconds(100)
});
Όταν ένα σήμα σε CancelOnSignals ανιχνεύεται, νέα στοιχεία παραλείπονται (μετρούνται ως αποτυχημένα).
Όταν ένα σήμα σε DeferOnSignals ανιχνεύεται, νέα αντικείμενα περιμένουν μέχρι να καθαρίσει το σήμα.
Από Adaptive TranslationService.cs:
public class AdaptiveTranslationService : IAsyncDisposable
{
private readonly EphemeralWorkCoordinator<TranslationRequest> _coordinator;
private readonly ITranslationApi _translationApi;
public AdaptiveTranslationService(ITranslationApi translationApi)
{
_translationApi = translationApi;
_coordinator = new EphemeralWorkCoordinator<TranslationRequest>(
ProcessTranslationAsync,
new EphemeralOptions
{
MaxConcurrency = 8,
MaxTrackedOperations = 100,
// New work is deferred while any "rate-limit" or "rate-limit:*" signal is present
DeferOnSignals = new HashSet<string> { "rate-limit", "rate-limit:*" },
MaxDeferAttempts = 10,
DeferCheckInterval = TimeSpan.FromMilliseconds(100)
});
}
public async Task TranslateAsync(TranslationRequest request)
{
// Optional: extra politeness based on most recent retry-after
var rateLimitSignals = _coordinator.GetSignalsByPattern("rate-limit:*");
if (rateLimitSignals.Count > 0)
{
var latest = rateLimitSignals
.OrderByDescending(s => s.Timestamp)
.First()
.Signal; // "rate-limit:5000ms"
if (TryParseRetryAfter(latest, out var delay))
{
await Task.Delay(delay);
}
}
await _coordinator.EnqueueAsync(request);
}
public static bool TryParseRetryAfter(string signal, out TimeSpan delay)
{
delay = default;
var parts = signal.Split(':', 2);
if (parts.Length != 2) return false;
var payload = parts[1].Trim();
if (!payload.EndsWith("ms", StringComparison.OrdinalIgnoreCase)) return false;
var numPart = payload[..^2];
if (!int.TryParse(numPart, out var ms) || ms < 0) return false;
delay = TimeSpan.FromMilliseconds(ms);
return true;
}
}
Κάθε περίπτωση αυτής της υπηρεσίας υποχωρεί αυτόματα όταν χτυπηθούν τα όρια των τιμών, δεν υπάρχει κοινή κατάσταση, δεν περνάει κανένα μήνυμα, απλά διαβάζω το εφήμερο παράθυρο.
Πολλαπλοί συντονιστές μπορούν να αισθανθούν ο ένας τον άλλον μέσω ενός κοινού SignalSink:
public class OrderProcessingSystem
{
private readonly SignalSink _sharedSignals = new(maxCapacity: 1000);
private readonly EphemeralWorkCoordinator<Order> _orderProcessor;
private readonly EphemeralWorkCoordinator<PaymentRequest> _paymentProcessor;
public OrderProcessingSystem()
{
var options = new EphemeralOptions { Signals = _sharedSignals };
_orderProcessor = new EphemeralWorkCoordinator<Order>(
ProcessOrderAsync, options);
_paymentProcessor = new EphemeralWorkCoordinator<PaymentRequest>(
ProcessPaymentAsync, options);
}
public async Task ProcessOrderAsync(Order order)
{
// Check shared signals for payment gateway issues
if (_sharedSignals.Detect("gateway-error"))
{
await _retryQueue.EnqueueAsync(order);
return;
}
await _orderProcessor.EnqueueAsync(order);
}
}
[HttpGet("/health/detailed")]
public IActionResult GetDetailedHealth()
{
return Ok(new
{
translation = new
{
pending = _translationCoordinator.PendingCount,
active = _translationCoordinator.ActiveCount,
recentRateLimits = _translationCoordinator.CountSignals("rate-limit"),
recentTimeouts = _translationCoordinator.CountSignals("timeout"),
recentSuccess = _translationCoordinator.CountSignals("success"),
hasErrors = _translationCoordinator.HasSignalMatching("error.*")
},
payment = new
{
pending = _paymentCoordinator.PendingCount,
gatewayErrors = _paymentCoordinator.CountSignals("gateway-error"),
declines = _paymentCoordinator.CountSignals("declined"),
approvals = _paymentCoordinator.CountSignals("approved")
}
});
}
Απλά ρώτησε το εφήμερο παράθυρο.
Από SignalBasedCircuitBreaker.cs:
public class SignalBasedCircuitBreaker
{
private readonly string _failureSignal;
private readonly int _threshold;
private readonly TimeSpan _windowSize;
public SignalBasedCircuitBreaker(
string failureSignal = "failure",
int threshold = 5,
TimeSpan? windowSize = null)
{
_failureSignal = failureSignal;
_threshold = threshold;
_windowSize = windowSize ?? TimeSpan.FromSeconds(30);
}
public bool IsOpen<T>(EphemeralWorkCoordinator<T> coordinator)
{
var recentFailures = coordinator.GetSignalsSince(
DateTimeOffset.UtcNow - _windowSize);
return recentFailures.Count(s => s.Signal == _failureSignal) >= _threshold;
}
public int GetFailureCount<T>(EphemeralWorkCoordinator<T> coordinator)
{
var recentFailures = coordinator.GetSignalsSince(
DateTimeOffset.UtcNow - _windowSize);
return recentFailures.Count(s => s.Signal == _failureSignal);
}
}
// Usage
var circuitBreaker = new SignalBasedCircuitBreaker("api-error", threshold: 3);
if (circuitBreaker.IsOpen(_coordinator))
{
throw new CircuitOpenException("Too many recent API errors");
}
await _coordinator.EnqueueAsync(request);
Ο διακόπτης κυκλώματος δεν έχει δική του κατάσταση - απλώς διαβάζει το εφήμερο παράθυρο.
Από Σήματα. cs:
Όταν τα σήματα μπορούν να προκαλέσουν άλλα σήματα, διακινδυνεύετε άπειρους βρόχους. SignalConstraints αποτρέπει αυτό:
var options = new EphemeralOptions
{
SignalConstraints = new SignalConstraints
{
// Max propagation depth before blocking
MaxDepth = 10,
// Prevent A → B → A cycles
BlockCycles = true,
// Signals that end propagation chains
TerminalSignals = new HashSet<string> { "completed", "failed", "resolved" },
// Signals that emit but don't propagate
LeafSignals = new HashSet<string> { "logged", "metric" },
// Callback when a signal is blocked
OnBlocked = (signal, reason) =>
{
_logger.LogWarning("Signal {Signal} blocked: {Reason}",
signal.Signal, reason);
}
}
};
Παρακολουθήστε την αιτιώδη συνάφεια με EmitCaused:
public void HandleSignal(SignalEvent evt, ISignalEmitter emitter)
{
if (evt.Is("order-placed"))
{
// This signal carries the propagation chain
// Will be blocked if it would create a cycle
emitter.EmitCaused("inventory-reserved", evt.Propagation);
}
}
Η αλυσίδα διάδοσης παρακολουθεί το μονοπάτι: order-placed → inventory-reserved → ...
Εάν inventory-reserved προσπάθησε να εκπέμψει order-placed, θα ήταν μπλοκαρισμένο (ανιχνεύθηκε κύκλος).
Από Σήματα. cs:
Για σήματα που πρέπει να είναι ορατά σε όλους τους συντονιστές:
public sealed class SignalSink
{
private readonly ConcurrentQueue<SignalEvent> _window;
private readonly int _maxCapacity;
private readonly TimeSpan _maxAge;
public SignalSink(int maxCapacity = 1000, TimeSpan? maxAge = null);
// Raise signals
public void Raise(SignalEvent signal);
public void Raise(string signal, string? key = null);
// Sense signals
public IReadOnlyList<SignalEvent> Sense();
public IReadOnlyList<SignalEvent> Sense(Func<SignalEvent, bool> predicate);
public bool Detect(string signalName);
public bool Detect(Func<SignalEvent, bool> predicate);
public int Count { get; }
}
Χρήση:
// Create a shared sink
var sink = new SignalSink(maxCapacity: 1000, maxAge: TimeSpan.FromMinutes(2));
// Configure coordinators to use it
var options = new EphemeralOptions { Signals = sink };
// Or raise signals directly
sink.Raise("system-maintenance");
// Sense from anywhere
if (sink.Detect("system-maintenance"))
{
await DeferWorkAsync();
}
Στυλ Glob που ταιριάζει με το φίλτρο σήματος:
// Exact match
coordinator.HasSignal("rate-limit");
// Wildcard patterns
coordinator.HasSignalMatching("http.*"); // http.timeout, http.error
coordinator.HasSignalMatching("error.*.critical"); // error.payment.critical
coordinator.HasSignalMatching("user-???-failed"); // user-123-failed
// Comma-separated patterns in CancelOnSignals/DeferOnSignals
new EphemeralOptions
{
CancelOnSignals = new HashSet<string>
{
"system-overload, circuit-open", // Either pattern
"error.*" // Any error signal
}
}
Κρατήστε τα σήματα απλά και συνεπή:
// Good - simple, categorical
op.Signal("success");
op.Signal("failure");
op.Signal("rate-limit");
op.Signal("timeout");
op.Signal("cache-hit");
// Good - structured for parsing
op.Signal("rate-limit:5000ms");
op.Signal("retry:attempt-3");
op.Signal("slow:2500ms");
op.Signal("http.error:429");
// Good - hierarchical for pattern matching
op.Signal("payment.declined");
op.Signal("payment.approved");
op.Signal("payment.gateway-error");
// Avoid - entity identification belongs in Key, not signals
op.Signal("user-123-rate-limited"); // Bad
// Instead
op.Key = "user-123";
op.Signal("rate-limit");
Χειριστές συγχρονικού σήματος (OnSignalΓια να κάνετε I/O-δεμένη εργασία, σήματα ανεμιστήρα σε μια διαδρομή async με SignalDispatcher (ταύτιση πίνακα, καθορισμένη σειρά) ή AsyncSignalProcessor.
await using var dispatcher = new SignalDispatcher(new EphemeralOptions
{
MaxConcurrency = Environment.ProcessorCount,
MaxConcurrencyPerKey = 1 // sequential per signal name by default
});
dispatcher.Register("error.*", evt => _alerts.SendAsync(evt.Signal));
dispatcher.Register("progress:*", evt => _metrics.Record(evt.Signal));
// In coordinator options, keep OnSignal chor options, keep OnSignal cheap and enqueue
var coordinator = new EphemeralWorkCoordinator<Request>(
body,
new EphemeralOptions
{
OnSignal = dispatcher.Dispatch
});
Υποστήριξη μοτίβα *, ?, και λίστες κόμμα ("error.*,timeout"). Όλοι οι χειριστές που ταιριάζουν λειτουργούν με εντολή εγγραφής σε έναν συντονιστή με κλειδί φόντο· η εκπομπή παραμένει συγχρονισμένη.
Για αυτόνομη επεξεργασία async:
await using var processor = new AsyncSignalProcessor(
async (signal, ct) =>
{
await _externalService.LogAsync(signal, ct);
},
maxConcurrency: 4,
maxQueueSize: 1000);
// Enqueue signals (returns immediately)
processor.Enqueue(new SignalEvent(
"rate-limit",
operationId,
key,
DateTimeOffset.UtcNow));
Ένα πλήρες παράδειγμα που συνδυάζει την επεξεργασία σήματος async με την ενσωμάτωση τηλεμετρίας:
public class TelemetrySignalHandler : IAsyncDisposable
{
private readonly AsyncSignalProcessor _processor;
private readonly ITelemetryClient _telemetry;
public TelemetrySignalHandler(ITelemetryClient telemetry)
{
_telemetry = telemetry;
_processor = new AsyncSignalProcessor(
HandleSignalAsync,
maxConcurrency: 8,
maxQueueSize: 5000);
}
// Synchronous entry point - returns immediately
public bool OnSignal(SignalEvent signal) => _processor.Enqueue(signal);
private async Task HandleSignalAsync(SignalEvent signal, CancellationToken ct)
{
var properties = new Dictionary<string, string>
{
["signal"] = signal.Signal,
["operationId"] = signal.OperationId.ToString(),
["key"] = signal.Key ?? "none"
};
await _telemetry.TrackEventAsync("EphemeralSignal", properties, ct);
// Categorized tracking based on signal prefix
if (signal.StartsWith("error"))
await _telemetry.TrackExceptionAsync(signal.Signal, properties, ct);
else if (signal.StartsWith("perf"))
await _telemetry.TrackMetricAsync(signal.Signal, 1, ct);
}
// Expose stats for monitoring
public int QueuedCount => _processor.QueuedCount;
public long ProcessedCount => _processor.ProcessedCount;
public long DroppedCount => _processor.DroppedCount;
public async ValueTask DisposeAsync() => await _processor.DisposeAsync();
}
Σύνδεσέ το στον συντονιστή σου:
await using var telemetryHandler = new TelemetrySignalHandler(telemetryClient);
await using var coordinator = new EphemeralWorkCoordinator<Request>(
ProcessAsync,
new EphemeralOptions
{
OnSignal = signal => telemetryHandler.OnSignal(signal)
});
Ο χειριστής:
OnSignal επιστρέφει αμέσωςΤα σήματα είναι ισχυρά επειδή είναι εφεμέραλCity name (optional, probably does not need a translation):
Το εφήμερο παράθυρο είναι ήδη εκεί για αποσφαλμάτωση.
Τα σήματα μετατρέπουν τα μεμονωμένα άτομα εκτέλεσης σε ένα δίκτυο ανίχνευσηςΚάθε συντονιστής μπορεί:
CancelOnSignals και DeferOnSignalsΧωρίς μεσίτη μηνυμάτων, χωρίς κοινόχρηστο κράτος, χωρίς πρωτόκολλο συντονισμού, μόνο με μεταδεδομένα που φυσικά αποσυντίθενται.
Τα άτομα δεν μιλάνε απευθείας μεταξύ τους - αφήνουν ίχνη στο εφήμερο παράθυρο που μπορούν να παρατηρήσουν οι άλλοι.
Φωτιά... σήμα... αίσθηση... ξέχασε.
© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.