Een klein primitief dat gelijktijdig werk verandert in een gecoördineerd, adaptief systeem.
"The Efemeral Signalen Pattern"
In Deel 1 Wij bouwden efemerale uitvoering - begrensd, prive, zelfreinigende async workflows. Deel 2 We veranderden dat in een herbruikbare bibliotheek met coördinatoren, gesleutelde pijpleidingen en DI integratie.
Dit artikel voegt een kleine functie die alles verandert: signalen.
Dit is nu in de meest lucid.efemerals Nuget pakket ook meer dan 20 meestal lucid.efemerals patronen en 'atomen'.
De volledige broncode staat in de meestallucid.atoms GitHub repository
De signaalinfrastructuur leeft in:
Bestand Doel
| ------ | --------- |
|---|---|
| EfemeralOperation.cs Signaalemissie en terugtrekking uit de activiteiten | |
EfemeralOptions.cs Signaalreactieve configuratie (CancelOnSignals, DeferOnSignals, OnSignal, OnSignalRetracted) |
|
| StringPatternMatcher.cs Glob-stijl patroon dat overeenkomt met het signaal filteren | |
SignalDispatcher.cs Async signaal routing met patroon matching (ondersteunt *, ?, kommalijsten, deterministische volgorde) |
|
| Voorbeelden/SignalingHttpClient.cs Sample fijnkorrelige signaalemissie voor HTTP-gesprekken | |
| Voorbeelden/AdaptiveTranslationService.cs Adaptive rate limiting with signal-based outillage | |
| Voorbeelden/SignalBasedCircuitBreaker.cs Een Circuit-onderbreker die het kortstondige signaalvenster leest | |
| Voorbeelden/TelemetrieSignalHandler.cs Async signaalverwerking met telemetrie-integratie |
Onze kortstondige coördinatoren zijn goed in het verwerken van werk, maar ze zijn geïsoleerd. Elke coördinator weet van zijn eigen activiteiten, maar heeft geen besef van wat er elders in het systeem gebeurt.
// 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
We kunnen expliciete afhankelijkheden verbinden, maar dat zorgt voor koppeling. omgevingsbewustzijn - coördinatoren die hun omgeving kunnen voelen zonder direct verbonden te zijn.
Signalen laten executoriale atomen sporen achter in hun kortstondige venster. Die sporen gedragen zich als kortstondige feiten:
Coördinators kunnen hun gedrag dan veranderen op basis van de signalen zichtbaar in hun venster. Dat is omgevingsbewustzijn zonder afhankelijkheden.
Van Signalen.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);
}
Een operatie kan signalen oproepen tijdens de uitvoering. Die signalen leven in het efemerale venster naast de operatie. Wanneer de operatie veroudert, gaan de signalen mee.
Geen berichtenmakelaar, geen aparte infrastructuur, alleen verplichtingen aan operaties.
Omdat signalen in het kortstondige venster leven, erven ze haar garanties: begrensde grootte, automatische veroudering en nul levenscyclus overhead.
Een signaal op zich veroorzaakt geen executie. Het registreert alleen een feit in het kortstondige venster. Niets loopt omdat er een signaal werd uitgezonden.
Als een coördinator een OnSignal handler definieert, draait hij synchroon wanneer een signaal wordt uitgezonden - maar alleen omdat de coördinator ervoor koos om het te bevestigen. Emitters weten het niet of geven er niks om. Het verwijderen van alle handlers laat het gedrag van de kern onveranderd.
Een signaal is alleen verbonden met het atoom/operatie dat het uitstraalde. Geen enkel signaal muteert of annoteert een ander atoom. Er is geen gedeelde beschrijfbare bus.
Een signaal negeren voegt een feit toe aan de geschiedenis van het atoom. Signalen worden nooit bijgewerkt of overschreven. Retractions verwijderen alleen de outlet eigen signalen.
Signalen bestaan alleen binnen het ephemerale venster van de coördinator. Ze verlopen automatisch als het venster veroudert. Er wordt niets volgehouden tenzij je expliciet doorzettingsvermogen opbouwt.
Wanneer een coördinator signalen controleert op sleutel K., scant hij:
de atomen in het venster
en de lokale signalen voor die atomen Waarnemers veranderen nooit een atoomsignaaltoestand.
Als u signalen wilt sturen naar async workflows, moet u het volgende gebruiken:
SignalDispatcher
AsyncSignalProcessor
of andere adapters.
Dit zijn optionele lagen gebouwd bovenop signalen, geen deel van hun semantiek.
Geen atoom mag de snapshot, toestand, signalen of metadata van een ander atoom veranderen. Coördinatie vindt plaats door:
signalen
detectie
vensters
beleid
Niet door schrijven.
Een SignalSink of een zwerm-gecompliceerd oppervlak mag een gecombineerd beeld van signalen tonen Maar het is altijd alleen-lezen, nooit gezaghebbend, nooit beschrijfbaar.
Als alle verwerkers (OnSignal, dispatchers, processors) zijn losgekoppeld, het systeem blijft volledig correct en voorspelbaar. Signalen hebben nog steeds betekenis omdat het feiten zijn, geen triggers.
Gebeurtenissen zijn als een telefoontje:
"Ik bel je nu, neem op en reageer."
Signalen zijn als voetafdrukken in de sneeuw:
"Ik liet voetafdrukken achter. Als je wilt weten waar ik heen ging - kijk. Als het je niet kan schelen - negeer het."
Dit is de reden waarom signalen nooit breken, nooit blokkeren, en nooit interactie met de controlestroom tenzij u kiezen om ze te peilen.
Het onderscheid is belangrijk omdat signalen elimineren:
Probleem Gebeurtenissen hebben het Signalen Vermijd het |---------|:--------------:|:----------------:| Koppeling Uitgever → Abonnees Geen Timing afhankelijkheden Onmiddelijke reactie vereist, poll als je klaar bent Feedback loops . . Handlers kunnen triggeren . . Geen lussen tenzij uitdrukkelijk gevraagd . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . Bestel semantiek Bestelzaken Irrelevant Bezorgingsgarantie moet geleverd/afgehandeld worden Geen levering, gewoon bestaan Fout bij het verspreiden van de handlerfouten . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . Reentrancy risks Common Impossible Een handler faalt, ketting breekt, geen ketting om te breken
Feature, Evenementen, Signalen, Signalen, Signalen, Signalen, Signalen, Signalen, Signalen, Signalen, Signalen, Signalen, Signalen, Signalen, Signalen, Signalen, Signalen, Signalen, Signalen, Signalen, Signalen, Signalen, Signalen, enz. |---------|--------|---------| Sterkte (uitgever → abonnees) Geen Timing Onmiddelijke ommezwaai Levering Gegarandeerd/verzorgd Geen levering, gewoon bestaan De reactie die nodig is optioneel Data Payload Vaak zwaar Tiny string metadata . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . Opbergen Geen Schuifbaar LRU-stijlvenster Leven en leven Instantane Verduisteringen automatisch Failure modi Veel bijna geen
Aanpak van gebeurtenissen (klassiek):
public event Action RateLimited;
try
{
await CallApiAsync();
}
catch (RateLimitException)
{
RateLimited?.Invoke(); // Makes someone else act right now
}
Problemen:
Signaalnadering (efemeraal):
try
{
await CallApiAsync(ct);
}
catch (RateLimitException ex)
{
op.Signal($"rate-limit:{ex.RetryAfterMs}ms");
throw;
}
Er gebeurt hier niks.
if (translator.HasSignal("rate-limit"))
{
await Task.Delay(1000); // Act when *we* choose
}
Opmerkingen:
Gebeurtenissen overdracht controle. Signalen overdracht context.
Dat is het hele mentale model in één zin.
Afhankelijk van uw achtergrond:
Definitie van het publiek |----------|------------| | Systeemdenkers Een stenig allergisch substraat voor indirecte coördinatie | Ingenieurs Lichtgewicht metadata verbonden aan handelingen in een begrensd schuifvenster | PL/Concurrency Nerds Een impliciete, temporale schoolbord gekoppeld aan begrensde concurrency semantics | Gebruikers van het kader Een in-proces, zelfreinigende staat oppervlak dat u kunt query op elk gewenst moment
Signalen zijn sporen achtergelaten in een gedeeld, begrensd geheugenoppervlak. Iedereen kan er naar kijken. Niemand is verplicht om te reageren. Dit is een fundamenteel Stigmereng model van coördinatie - dezelfde mieren gebruiken, dezelfde schoolbord systemen gebruikt in de vroege AI, en dezelfde moderne CRDT roddel netwerken hint op.
using var activity = source.StartActivity("ProcessOrder");
activity?.SetTag("order.id", orderId);
activity?.SetTag("rate.limited", true);
Beste voor: Gedistribueerde tracering over diensten, lange termijn telemetrie-opslag, correlatie-ID's.
Telemetrie gebruiken wanneer: U moet verzoeken traceren over meerdere diensten, metrics opslaan voor analyse, of integreren met monitoring tools.
Ephemerale signalen gebruiken wanneer: Je hebt in-proces omgevingsbewustzijn nodig, reactieve coördinatie, of wil geen telemetrie infrastructuur.
var rateLimits = Observable.FromEventPattern<RateLimitEventArgs>(
h => api.RateLimitHit += h,
h => api.RateLimitHit -= h);
rateLimits
.Throttle(TimeSpan.FromSeconds(1))
.Subscribe(e => HandleRateLimit(e));
Beste voor: Complex event processing, time-based operations, combining multiple event streams.
Rx gebruiken wanneer: Je hebt complexe tijdsvragen nodig (vensters, debouncing, het combineren van stromen).
Ephemerale signalen gebruiken wanneer: U wilt eenvoudiger peiling gebaseerde sensing, automatische opruiming, of integratie met operatie tracking.
public class RateLimitNotification : INotification
{
public int RetryAfterMs { get; init; }
}
await _mediator.Publish(new RateLimitNotification { RetryAfterMs = 5000 });
Beste voor: Ontkoppelde in-proces event handling met meerdere handlers.
MediatR gebruiken wanneer: Je wilt dat meerdere handlers synchroon reageren op dezelfde gebeurtenis.
Ephemerale signalen gebruiken wanneer: U wilt ambient sensing zonder expliciete abonnement, zelfreinigende geschiedenis, of integratie met begrensde uitvoering.
var circuitBreaker = Policy
.Handle<HttpRequestException>()
.CircuitBreakerAsync(5, TimeSpan.FromSeconds(30));
Beste voor: Veerkracht rond individuele gesprekken met automatisch staatsbeheer.
Polly gebruiken wanneer: Je hebt per-call veerkracht nodig met automatische half-open/gesloten overgangen.
Ephemerale signalen gebruiken wanneer: U wilt ambient awareness over vele operaties, aangepaste circuit logica, of integratie met operatie tracking.
Combineer ze: Gebruik Polly in je werklichaam, zend signalen uit wanneer circuits struikelen.
Approach "Zelfreinigend "Ambtenaren Sensing "Ontkoppeld "Aangepaste Logica "Integratie " |----------|:-------------:|:---------------:|:---------:|:------------:|:-----------:| OpenTelemetry OpenTelemetrie Uitwendig gereedschap Reactieve uitbreidingen . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . MediatR, ❌, ❌, Polly Circuit Breaker . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . | Ephemeral Signalen Ingebouwd ingebouwd ingebouwd ingebouwd
Uitvoering van concrete acties 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; }
}
Binnenin uw werklichaam:
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;
}
});
Signalen zijn slechts tekenreeksen. Gebruik eenvoudige namen ("rate-limit") of gestructureerde namen ("rate-limit:5000ms").
Patroonfilters gebruiken glob semantiek (*, ?) en ondersteuning van kommalijsten ("error.*,timeout") Matching is deterministisch en allocatie-licht via StringPatternMatcher.
Voor zeer gedetailleerde opmerkzaamheid, kunt u signalen uit te zenden in elke fase van een operatie. De bibliotheek bevat een sample SignaliserenHttpClient dat dit patroon laat zien:
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...
});
Dit zendt signalen uit in elk stadium:
Signaal Wanneer:
|--------|------|
| stage.starting Voordat het verzoek begint
| progress:0 Initiële vooruitgang
| stage.request Verzonden HTTP-verzoek
| stage.headers Reactie headers ontvangen
| stage.reading Beginnen met het lezen van lichaam
| progress:XX Voortgangspercentage (0-100) tijdens het downloaden
| stage.completed Download voltooid
U kunt deze vervolgens opvragen met een patroon dat overeenkomt met:
// 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
}
Operaties kunnen ook hun eigen signalen verwijderen. Dit is nuttig voor tijdelijke staten:
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;
}
});
Net als signaalemissie, kunnen intrekkingen terugroepen:
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);
}
});
De SignalRetractedEvent omvat:
Signal - De ingetrokken signaalnaamOperationId - De operatie die het introk.Key - De sleutel van de operatie (indien aanwezig)Timestamp - Wanneer het terugtrekken plaatsvondWasPatternMatch - Waar indien ingetrokken via RetractMatchingPattern - Het gebruikte patroon (als het patroon overeenkomt)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;
}
});
Alle coördinatoren bieden geoptimaliseerde signaalquering:
// 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");
Van EfemeralOptions.cs:
Coördinatoren kunnen automatisch reageren op signalen:
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)
});
Wanneer een signaal binnenkomt CancelOnSignals wordt gedetecteerd, nieuwe items worden overgeslagen (geteld als mislukt).
Wanneer een signaal binnenkomt DeferOnSignals wordt gedetecteerd, nieuwe items wachten tot het signaal weg is.
Van AdaptiveTranslationService.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;
}
}
Elke instantie van deze dienst gaat automatisch achteruit wanneer de snelheidslimieten worden geraakt. Geen gedeelde status. Geen bericht passeren. Gewoon het ephemerale venster lezen.
Meerdere coördinatoren kunnen elkaar voelen via een gedeelde 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")
}
});
}
Geen metrics-bibliotheek nodig. Vraag gewoon naar het efemerale venster.
Van 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);
De stroomonderbreker heeft geen eigen staat - hij leest alleen het kortstondige venster.
Van Signalen.cs:
Wanneer signalen andere signalen kunnen veroorzaken, riskeer je oneindige loops. SignalConstraints voorkomt dit:
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);
}
}
};
Track causaliteit met 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);
}
}
De voortplantingsketen volgt het pad: order-placed → inventory-reserved → ...
Als inventory-reserved probeerde uit te zenden order-placed, het zou worden geblokkeerd (cyclus gedetecteerd).
Van Signalen.cs:
Voor signalen die zichtbaar moeten zijn tussen coördinatoren:
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; }
}
Gebruik:
// 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-stijl matching voor signaalfiltering:
// 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
}
}
Houd signalen eenvoudig en consistent:
// 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");
Synchroonsignaalverwerkers (OnSignal) draaien op de thread van de operatie te houden ze snel. Om te doen I/O-gebonden werk, ventilator signalen in een async pad met SignalDispatcher (pattern matching, deterministische volgorde) of 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
});
Patronenondersteuning *, ?, en komma's ("error.*,timeout"Alle bijpassende handlers draaien in registratie volgorde op een achtergrond getoetste coördinator; emit blijft synchroon.
Voor standalone async-verwerking:
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));
Een compleet voorbeeld dat async signaalverwerking combineert met telemetrie-integratie:
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();
}
Sluit het aan op uw coördinator:
await using var telemetryHandler = new TelemetrySignalHandler(telemetryClient);
await using var coordinator = new EphemeralWorkCoordinator<Request>(
ProcessAsync,
new EphemeralOptions
{
OnSignal = signal => telemetryHandler.OnSignal(signal)
});
De begeleider:
OnSignal komt onmiddellijk terugSignalen zijn krachtig omdat ze kortstondig:
Onroerend goed . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . |----------|---------| | Vastgebonden Kan niet ongebonden groeien - oude signalen verouderen | Zelfreinigend Geen opschoningscode nodig | Ontkoppeld Emitters weten niets van luisteraars. | Waarneembaar Elke code kan de huidige toestand voelen. | Privé Geen gebruikersgegevens - alleen signaalnamen | Snel O(1) detectie met kortsluiting
Het ephemerale venster is er al voor debuggen. Signalen geven het semantische betekenis.
Signalen veranderen geïsoleerde uitvoeringsatomen in een sensornetwerkElke coördinator kan:
CancelOnSignals en DeferOnSignalsGeen berichtenmakelaar, geen gedeelde staat, geen coördinatieprotocol... alleen operaties met metadata die natuurlijk vervallen.
De atomen praten niet rechtstreeks met elkaar - ze laten sporen achter in het kortstondige venster dat anderen kunnen observeren.
Vuur... signaal... gevoel... vergeet het.
© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.