En liten primitiv som förvandlar samtidigt arbete till ett samordnat, adaptivt system.
"Det efemära signalmönstret"
Till Häfte 1 Vi byggde efemeral exekvering - avgränsade, privata, självrengörande async arbetsflöden. Häfte 2 Vi gjorde det till ett återanvändbart bibliotek med koordinatorer, nyckelrörledningar och DI-integration.
Denna artikel lägger till en liten funktion som ändrar allt: signaler.
Detta är nu också i det mest lucid.ephemerals Nuget-paketet Mer än 20 oftast lucid.ephemerala mönster och 'atomer'.
Den fullständiga källkoden finns i mestadels lucid.atomer GitHub arkiv
Signalinfrastrukturen lever i:
Filens syfte
| ------ | --------- |
|---|---|
| EfemeralOperation.cs till Signalutsläpp och reträtt från drift | |
Efemeralalternativ.cs Signalreaktiva konfigurationer (CancelOnSignals, DeferOnSignals, OnSignal, OnSignalRetracted) |
|
| StringPatternMatcher.cs på ett mönster som passar för signalfiltrering | |
Signalsändare.cs på Async signal routing med mönstermatchning (stöder *, ?, kommateckenlistor, deterministisk ordning) |
|
| Exempel/SignalingHttpClient.cs på prov finkornig signalemission för HTTP-samtal | |
| Exempel/AdaptiveTranslationService.cs Adaptive rate begränsande med signalbaserat uppskjutande | |
| Exempel/SignalBasedCircuitBreaker.cs till Circuit breaker avläsning efemeral signalfönster | |
| Exempel/TelemetriSignalHandler.cs till Async-signalbehandling med telemetriintegration |
Våra efemära samordnare är bra på att bearbeta arbete, men de är isolerade. Varje samordnare känner till sin egen verksamhet, men har ingen medvetenhet om vad som händer någon annanstans i systemet.
// 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
Vi kan koppla upp tydliga beroenden, men det skapar koppling. Omgivningsmedvetenhet - koordinatorer som kan känna sin omgivning utan att vara direkt anslutna.
Signaler låter avrättning atomer lämna spår i deras efemerala fönster. Dessa spår beter sig som kortlivade fakta:
Koordinatorerna kan då ändra sitt beteende baserat på de signaler som syns i deras fönster. Det är omgivningsmedvetenhet utan beroenden.
Från Signaler.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);
}
En operation kan höja signaler under utförandet. Dessa signaler lever i det efemära fönstret bredvid operationen. När operationen åldras ut, signalerna går med den.
Ingen meddelandemäklare, ingen separat infrastruktur, bara band kopplade till verksamheten.
Eftersom signalerna lever inne i det efemära fönstret ärver de sina garantier: begränsad storlek, automatiskt åldrande och noll livscykel overhead.
En signal, i sig själv, orsakar inte avrättning. Det finns bara ett faktum i det efemära fönstret. Inget går för att en signal sändes ut.
Om en koordinator definierar en OnSignal-hanterare, kör den synkront när en signal avges - men bara för att koordinatorn valde att ansluta den. Emitters bryr sig inte. Att ta bort alla hanterare lämnar kärnans beteende oförändrat.
En signal är endast fäst vid den atom/operation som avger den. Ingen signal någonsin mutera eller kommentera en annan atom. Det finns ingen delad skrivbar buss.
Att sända en signal lägger till ett faktum till atomens historia. Signaler uppdateras aldrig eller skrivs över. Retraktioner tar bara bort emitterns egna signaler.
Signaler finns bara inom samordnarens efemära fönster. De löper ut automatiskt när fönstret åldras ut. Inget finns kvar om du inte uttryckligen bygger upp ihärdighet.
När en koordinator kollar signaler för nyckel K på skannar den:
atomerna i dess fönster
och signalerna lokala till dessa atomer Observatörer ändrar aldrig ett atomtillstånd.
Om du vill ha signaler för att köra async-arbetsflöden måste du använda:
Signalavsändare
AsyncSignalProcessorName
eller andra adaptrar.
Dessa är valfria lager som byggs ovanpå signaler, inte en del av deras semantik.
Ingen atom får ändra en annan atoms ögonblicksbild, tillstånd, signaler eller metadata. Samordning sker genom
signaler
Känslor
fönster
politik
Inte genom att skriva.
En signalSänka eller svärmaggregaterad yta kan visa en kombinerad vy av signaler — Men den är alltid skrivbar, aldrig auktoritativ, aldrig skrivbar.
Om alla hanterare (OnSignal, Avsändare, Processorer) är fristående, systemet förblir helt korrekt och förutsägbart. Signaler har fortfarande betydelse eftersom de är fakta, inte utlösare.
Evenemang är som ett telefonsamtal:
"Jag ringer dig nu och svarar."
Signaler är som fotavtryck i snön:
"Jag lämnade fotavtryck, om du vill veta vart jag tog vägen, titta om du inte bryr dig, ignorera det."
Det är därför signaler aldrig går sönder, aldrig blockerar, och aldrig interagerar med kontrollflödet om du inte välj För att höra dem.
Skillnaden är viktig eftersom signaler eliminerar:
på problem på händelser har det med signaler undvika det på |---------|:--------------:|:----------------:| på Coupling på Utgivare → Prenumeranter och ingen på tid beroenden med omedelbar reaktion krävs på Ambient, opinionsundersökning när redo . Inga loopar om inte uttryckligen begärts Ordning av semantik och ordningsfrågor på Leveransgarantier på Måste levereras/handleds på Ingen leverans, bara existens på fel förökning på Handler fel propagera på Isolerad på "Reentreprenadrisker" Gemensamma "Omöjliga" En hanterare misslyckas, kedjan går sönder Ingen kedja att bryta
. Feature med händelser och signaler |---------|--------|---------| Kopiering på ett kraftfullt sätt (förlag → prenumeranter) Tidpunkt med omedelbar omskolning Ingen leverans, bara existens Obligatorisk reaktion på grund av att den är frivillig på data nyttolast och ofta tunga, små strängmetadata LRU-format fönster Livstid på ett ögonblick Deceays automatiskt på grund av fel, många på nästan inget sätt
Tillvägagångssätt vid händelser (klassiskt):
public event Action RateLimited;
try
{
await CallApiAsync();
}
catch (RateLimitException)
{
RateLimited?.Invoke(); // Makes someone else act right now
}
Problem:
Signalinflygning (ephemeral):
try
{
await CallApiAsync(ct);
}
catch (RateLimitException ex)
{
op.Signal($"rate-limit:{ex.RetryAfterMs}ms");
throw;
}
Inget händer här. Senare, någonstans helt annorlunda:
if (translator.HasSignal("rate-limit"))
{
await Task.Delay(1000); // Act when *we* choose
}
Anmärkningar:
Överföringskontroll av händelser, signalöverföringskontext.
Det är hela mentalmodellen i en mening.
Beroende på din bakgrund:
"Publiken" och "definitionen" |----------|------------| | Systemtänkare Ettstigmergiskt substrat för indirekt samordning | Ingenjörer till Lättviktsmetadata som är kopplade till verksamheten i ett avgränsat skjutfönster | PL/konkursnördar En implicit, temporal svart tavla parad med avgränsad konvergens semantik | Ramanvändare En självrengörande yta under processen som du kan fråga när som helst
Signaler finns kvar på en delad, avgränsad minnesyta. Vem som helst kan titta på dem. Ingen är skyldig att reagera. Sprängämne samordningsmodell - samma en myror använder, samma en svarta tavlan system som används i tidig AI, och samma en modern CRDT skvaller nätverk tips på.
using var activity = source.StartActivity("ProcessOrder");
activity?.SetTag("order.id", orderId);
activity?.SetTag("rate.limited", true);
Bäst för: Distribuerad spårning över tjänster, långsiktig telemetrilagring, korrelations-ID.
Använd telemetri när: Du måste spåra förfrågningar över flera tjänster, lagra mätvärden för analys, eller integrera med övervakningsverktyg.
Använd efemerala signaler när: Du behöver i processen omgivande medvetenhet, reaktiv samordning, eller vill inte telemetri infrastruktur.
var rateLimits = Observable.FromEventPattern<RateLimitEventArgs>(
h => api.RateLimitHit += h,
h => api.RateLimitHit -= h);
rateLimits
.Throttle(TimeSpan.FromSeconds(1))
.Subscribe(e => HandleRateLimit(e));
Bäst för: Komplex händelsebehandling, tidsbaserad verksamhet, som kombinerar flera händelseströmmar.
Använd Rx när: Du behöver komplexa temporal frågor (fönster, debouncing, kombinera strömmar).
Använd efemerala signaler när: Du vill ha enklare röstningsbaserad analys, automatisk rensning, eller integration med drift spårning.
public class RateLimitNotification : INotification
{
public int RetryAfterMs { get; init; }
}
await _mediator.Publish(new RateLimitNotification { RetryAfterMs = 5000 });
Bäst för: Frikopplad händelsehantering under processen med flera hanterare.
Använd MediatR när: Du vill att flera hanterare ska reagera på samma händelse synkront.
Använd efemerala signaler när: Du vill ha omgivningsanalys utan explicit prenumeration, självrengörande historia, eller integration med begränsad utförande.
var circuitBreaker = Policy
.Handle<HttpRequestException>()
.CircuitBreakerAsync(5, TimeSpan.FromSeconds(30));
Bäst för: Resiliens kring enskilda samtal med automatisk statlig förvaltning.
Använd Polly när: Du behöver per-samtalsresiliens med automatiska halvöppna/stängda övergångar.
Använd efemerala signaler när: Du vill ha en omgivningsmedvetenhet över många operationer, anpassad kretslogik, eller integration med operationsspårning.
Kombinera dem: Använd Polly inuti din arbetskropp, avge signaler när kretsar resa.
Tillvägagångssätt och självrengörande på ett sätt som gör det möjligt för människor att känna igen sig och koppla loss sig från sin förmåga att integrera sig |----------|:-------------:|:---------------:|:---------:|:------------:|:-----------:| på OpenTeletry . . Externa verktyg Reaktiva förlängningar . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . MediantR på engelska på engelska på engelska på svenska på svenska på svenska på svenska på svenska på svenska på svenska . på Polly Circuit Breaker på och vis . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . | Efemära signaler Obligatoriska uppgifter som krävs för att uppfylla de krav som anges i bilaga II till förordning (EU) nr 600/2014.
Genomförande av insatser 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; }
}
Inuti din arbetskropp:
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;
}
});
Signaler är bara strängar. Använd enkla namn ("rate-limit") eller strukturerade namn ("rate-limit:5000ms").
Mönsterfilter använder globsemantik (*, ?) och stödja kommateckenlistor ("error.*,timeout". Matchning är deterministisk och allokering-ljus via StringPatternMatcher.
För mycket detaljerad observerbarhet kan du sända ut signaler i varje steg av en operation. Biblioteket innehåller ett prov SignaleringHttpClient som visar detta mönster:
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...
});
Detta avger signaler i varje steg:
När det är dags för signalen
|--------|------|
| stage.starting Innan begäran börjar
| progress:0 Initialt förloppsmärke
| stage.request Skickad HTTP- begäran
| stage.headers Svarsrubriker som mottagits
| stage.reading Börjar läsa kroppen
| progress:XX Förloppsprocent (0-100) under nedladdning
| stage.completed Ladda ner färdig
Du kan sedan fråga dessa med mönster som matchar:
// 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
}
Verksamheten kan också ta bort sina egna signaler. Detta är användbart för tillfälliga tillstånd:
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;
}
});
Precis som signalutsläpp, kan indragningar utlösa återuppringningar:
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);
}
});
och SignalRetractedEvent omfattar följande:
Signal - Det indragna signalnamnetOperationId - Operationen som drog tillbaka den.Key - Operationens nyckel (om sådan finns)Timestamp - När reträtt inträffadeWasPatternMatch - Sant om de dras in via RetractMatchingPattern - Mönstret som används (om mönstret matchar)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;
}
});
Alla samordnare tillhandahåller optimerad signalförfrågan:
// 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");
Från Efemeralalternativ.cs:
Samordnare kan automatiskt reagera på signaler:
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)
});
När en signal kommer in CancelOnSignals upptäcks, nya objekt hoppar över (räknas som misslyckades).
När en signal kommer in DeferOnSignals upptäcks, nya objekt vänta tills signalen klarnar.
Från 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;
}
}
Varje instans av den här tjänsten backar automatiskt bort när hastighetsgränser träffas. Inget delat tillstånd. Inget meddelande passerar. Läser bara det efemära fönstret.
Flera samordnare kan känna varandra genom en delad 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")
}
});
}
Det behövs inget databibliotek, fråga bara det efemära fönstret.
Från 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);
Kretsbrytaren har inget eget tillstånd - den läser bara det efemära fönstret.
Från Signaler.cs:
När signaler kan orsaka andra signaler, riskerar du oändliga slingor. SignalConstraints förhindrar detta:
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);
}
}
};
Spårets kausalitet med 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);
}
}
Utbredningskedjan spårar stigen: order-placed → inventory-reserved → ...
Om inventory-reserved försökte släppa ut order-placed, Det skulle blockeras (cykel upptäcks).
Från Signaler.cs:
För signaler som måste vara synliga mellan samordnare:
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; }
}
Användning:
// 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();
}
Från StringPatternMatcher.cs:
Matchning i Glob-stil för signalfiltrering:
// 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
}
}
Håll signalerna enkla och konsekventa:
// 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");
Synkrona signalhanterare (OnSignal) kör på operationens tråd – håll dem snabbt. För att utföra I/O-bundet arbete, signalerar fläkten in i en async-väg med SignalDispatcher (mönstermatchning, deterministisk ordning) eller 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
});
Stöd för mönster *, ?, och kommalistor ("error.*,timeout"). Alla matchande hanterare körs i registreringsordning på en bakgrundsnyckel koordinator; utsläpp förblir synkroniserade.
För fristående async-behandling:
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));
Ett komplett exempel som kombinerar async-signalbehandling med telemetriintegration:
Källa: TelemetriSignalHandler.cs
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();
}
Koppla upp det till din koordinator:
await using var telemetryHandler = new TelemetrySignalHandler(telemetryClient);
await using var coordinator = new EphemeralWorkCoordinator<Request>(
ProcessAsync,
new EphemeralOptions
{
OnSignal = signal => telemetryHandler.OnSignal(signal)
});
Handläggaren:
OnSignal återvänder omedelbartSignaler är kraftfulla eftersom de är Förgängliga:
Förmåner för egendom |----------|---------| | Gränsad Kan inte växa obundna - gamla signaler åldras ut | Självrengörande Ingen rensningskod behövs | Frikopplat "Emitters vet inte om lyssnare" | Observerbar Alla koder kan känna av det nuvarande tillståndet | Privat Inga användardata - bara signalnamn | Snabbt O(1) detektering med kortslutning
Efemerala fönstret är redan där för felsökning. Signaler ger det bara semantisk mening.
Signaler förvandlar isolerade avrättningsatomer till en Avkänning av nätverk. Varje samordnare kan:
CancelOnSignals och DeferOnSignalsIngen meddelandemäklare, inget delat tillstånd, inget samordningsprotokoll, bara operationer med metadata som naturligt sönderfaller.
Atomerna pratar inte direkt med varandra - de lämnar bara spår i det efemära fönstret som andra kan observera. Det ärstigmergy för async-system.
Eld, signal, sinne, glöm det.
© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.