Pieni alkukantainen, joka muuttaa samanaikaisen työn koordinoiduksi, mukautuvaksi järjestelmäksi.
"Efeeristen signaalien malli"
Sisään Osa 1 Rakensimme tilapäistä teloitusta - rajallisia, yksityisiä, itsepuhdistavia async-työvirtoja. 2 osa Teimme siitä uudelleenkäytettävän kirjaston, jossa on koordinaattoreita, keulaputkistoja ja DI-integraatio.
Tässä artikkelissa lisätään yksi pieni piirre, joka muuttaa kaiken: signaaleja.
Tämä on nyt enimmäkseen lucid.efemerals Nuget-paketti myös yli 20 enimmäkseen lucid.efemeral-kuviota ja 'atoms'.
Koko lähdekoodi on enimmäkseen lucid.atoms GitHub-varasto
Signaaliinfrastruktuuri elää:
Tiedoston tarkoitus
| ------ | --------- |
|---|---|
| EphemeralOperation.cs Signaalien päästöt ja vetäytyminen toiminnasta | |
EphemeralOptions.cs Signaalin reagoiva konfiguraatio (CancelOnSignals, DeferOnSignals, OnSignal, OnSignalRetracted) |
|
| StringPatternMatcher.cs Glob-tyylinen malli, joka vastaa signaalin suodatusta | |
SignalDispatcher.cs Async-signaalin reititys kuvioiden kanssa (tukia) *, ?, pilkkulistat, deterministinen järjestys) |
|
| Esimerkkejä/SignalingHttpClient.cs Näytteen hienorakeinen signaaliemissio HTTP-puheluissa | |
| Esimerkkejä/käännöspalvelu.cs Adaptiivisen koron rajoitus signaaliperusteisella lykkäyksellä | |
| Esimerkkejä/SignalBasedCircuitBreaker.cs Virtapiirin katkaisimen lukema hetkellinen signaaliikkuna | |
| Esimerkkejä/telemetriaSignalHandler.cs Async-signaalinkäsittely ja telemetriaintegraatio |
Evoluutiokoordinaattorimme ovat hyviä prosessoinnissa, mutta he ovat eristyksissä. Jokainen koordinaattori tietää omasta toiminnastaan, mutta ei tiedä, mitä muualla järjestelmässä tapahtuu.
// 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
Voisimme kytkeä yhteen selvät riippuvuudet, mutta se luo yhteyden. ympäristötietoisuus - koordinaattoreita, jotka voivat aistia ympäristönsä ilman suoraa yhteyttä.
Signaalit antavat teloitusatomien jättää jälkiä hetkelliseen ikkunaansa. Nämä jäljet käyttäytyvät kuin lyhytikäiset faktat:
Koordinaattorit voivat sitten muuttaa käytöstään ikkunassa näkyvien signaalien perusteella. Se on ympäristön tiedostamista ilman riippuvuuksia.
From Signaalit.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);
}
Operaatio voi nostaa signaaleja suorituksen aikana. Signaalit elävät operaation ohessa hetkellisessä ikkunassa. Kun operaatio vanhenee, signaalit kulkevat mukana.
Siinä kaikki, ei viestinvälittäjää, ei erillistä infrastruktuuria, vain toimintaan liittyviä ehtoja.
Koska signaalit elävät hetkellisen ikkunan sisällä, ne perivät sen takuut: rajatun koon, automaattisen vanhenemisen ja nollan elinkaaren yläpuolella.
Signaali itsessään ei aiheuta teloitusta. Se tallentaa vain faktan hetkellisessä ikkunassa. Mikään ei kulje, koska signaali lähetettiin.
Jos koordinaattori määrittelee OnSignal-käsittelijän, se toimii synkronisesti signaalin lähettäessä, mutta vain siksi, että koordinaattori päätti liittää sen. Lähettäjät eivät tiedä eivätkä välitä. Kaikkien käsittelijöiden poistaminen jättää ydinkäyttäytymisen ennalleen.
Signaali kiinnittyy vain sen säteilevään atomiin/toimintoon. Mikään signaali ei koskaan mutaatioita tai anota toista atomia. Ei ole yhteistä kirjoitettavaa bussia.
Signaalin lähettäminen lisää tosiasian atomien historiaan. Signaaleja ei koskaan päivitetä tai ylikirjoiteta. Peruutukset poistavat vain lähettimien omat signaalit.
Signaaleja on vain koordinaattorin hetkellisessä ikkunassa. Ne vanhenevat automaattisesti, kun ikkuna vanhenee. Mitään ei itsepintaisesti tehdä, jos sinnikkyyttä ei nimenomaan rakenneta.
Kun koordinaattori tarkistaa Kööpenhamina-avaimen signaalit, se skannataan:
atomit ikkunassaan
ja paikalliset signaalit noille atomeille Tarkkailijat eivät koskaan muokkaa atomien signaalitilaa.
Jos haluat signaalien ohjaavan async-työnkulkua, sinun on käytettävä:
SignalDispator
AsyncSignalProcessor
tai muita adaptereita.
Nämä ovat valinnaisia kerroksia, jotka on rakennettu signaalien päälle, eivät osa niiden semantiikkaa.
Mikään atomi ei saa muuttaa toisen atomin kuvaa, tilaa, signaaleja tai metadataa. Yhteensovittaminen tapahtuu seuraavasti:
signaaleja
aistiminen
ikkunat
politiikka
Ei kirjoitusten kautta.
SignalSinkin tai parven erittelemällä pinnalla voi olla yhdistetty kuva signaaleista – Mutta se on aina vain luettavaa, ei koskaan arvovaltaista, ei koskaan kirjoitettavaa.
Jos kaikki käsittelijät (OnSignal, lähettäjät, prosessorit) irrotetaan, järjestelmä pysyy täysin oikeana ja ennustettavana. Signaaleilla on yhä merkitystä, koska ne ovat tosiasioita, eivät laukaisevia.
Tapahtumat ovat kuin puhelinsoitto:
"Soitan heti, vastaa ja reagoi."
Signaalit ovat kuin jalanjäljet lumessa:
"Jätin jalanjäljet. Jos haluat tietää, minne menin - katso, jos et välitä, jätä se huomiotta."
Tämän vuoksi signaalit eivät koskaan katkea, eivät koskaan estä eivätkä ole vuorovaikutuksessa ohjausvirran kanssa, ellet valitse mielipidemittaukseen.
Erolla on merkitystä, koska signaalit poistavat
Ongelmia Tapahtumat On Se Signals Vältä sitä |---------|:--------------:|:----------------:| Julkaisija Kytkentä Julkaisija → Ei tilaajia Välitön reagointi tarpeen Ympäristö, kysely, kun se on valmis Palautesilmukat Handlers voi käynnistää käsittelijät Ei silmukoita ellei nimenomaisesti kysytä Semantiikan tilaaminen Järjestysasiat Asiaankuulumattomat Toimitustakuut On toimitettava/käsiteltävä Ei toimituksia, vain olemassa Virheiden lisääntyminen Handler-virheiden lisääntyminen Yksinäisten virheiden lisääntyminen "Vaarallisuusvaarat" "Mahdotonta" Yksi käsittelijä epäonnistuu, ketju katkeaa Ei ketjua katkaista
Ominaisuus Tapahtumat Signaalit |---------|--------|---------| Strong Coupling (julkaisija → tilaajat) None (ei tilaajia) Välittömän toiminnan ajankohta Toimitus Taattu/suunniteltu Ei toimitusta, vain olemassa "Vapaaehtoinen reaktio" Datan hyötykuorma Usein raskas pikkujousen metadata Storage None Sliding LRU-tyylinen ikkuna Elämänaika Välittömästi tuhoutuu automaattisesti Hylkäystiloja Monia Ei juuri yhtään
Tapahtumia koskeva lähestymistapa (klassinen):
public event Action RateLimited;
try
{
await CallApiAsync();
}
catch (RateLimitException)
{
RateLimited?.Invoke(); // Makes someone else act right now
}
Ongelmia:
Signaalilähestyminen (hetkellinen):
try
{
await CallApiAsync(ct);
}
catch (RateLimitException ex)
{
op.Signal($"rate-limit:{ex.RetryAfterMs}ms");
throw;
}
Täällä ei tapahdu mitään, mutta myöhemmin, jossain aivan erilaisessa:
if (translator.HasSignal("rate-limit"))
{
await Task.Delay(1000); // Act when *we* choose
}
Huomautukset:
Tapahtumasiirron ohjaus, viestien siirtokonteksti.
Siinä on koko henkinen malli yhdessä lauseessa.
Riippuen taustastasi:
Yleisön määritelmä |----------|------------| | Järjestelmä-ajattelijat Epäsuoran koordinaation stimuloiva substraatti | Insinöörit Kevyet metatiedot, jotka on liitetty toimintaan suljetussa liukuikkunassa | PL/Concurrency Nords Implisiittinen, ajallinen liitutaulu, johon on yhdistetty rajoitettua sekavaluuttasemantiikkaa | Kehyskäyttäjät Prosessinaikainen, itsepuhdistava tilapinta, jota voi tiedustella milloin tahansa
Signaalit ovat jälkiä, jotka jäävät jaettuun, rajattuun muistipintaan. Kuka tahansa voi katsoa niitä. Kukaan ei ole velvollinen reagoimaan. stigmerginen Malli koordinaatiosta - samaa muurahaista käytetään, samaa yhtä liitutaulujärjestelmää käytetään tekoälyssä, ja samaa, johon modernit CRDT-juoruverkot vihjaavat.
using var activity = source.StartActivity("ProcessOrder");
activity?.SetTag("order.id", orderId);
activity?.SetTag("rate.limited", true);
Paras: Hajautettu jäljitys eri palveluihin, pitkän aikavälin telemetriatallennus, korrelaatiotunnukset.
Käytä telemetriaa, kun: Sinun täytyy jäljittää pyyntöjä useissa palveluissa, tallentaa mittarit analysointia varten tai integroitua seurantatyökaluihin.
Käytä efeerisiä signaaleja, kun: Tarvitset prosessinaikaista ympäristötietoisuutta, reaktiivista koordinaatiota tai et halua telemetriainfrastruktuuria.
var rateLimits = Observable.FromEventPattern<RateLimitEventArgs>(
h => api.RateLimitHit += h,
h => api.RateLimitHit -= h);
rateLimits
.Throttle(TimeSpan.FromSeconds(1))
.Subscribe(e => HandleRateLimit(e));
Paras: Monimutkaista tapahtumakäsittelyä, aikalähtöistä toimintaa, jossa yhdistetään useita tapahtumavirtoja.
Käytä Rx:ää, kun: Tarvitset monimutkaisia ajallisia kyselyitä (ikkunoita, nujertamista, virtojen yhdistämistä).
Käytä efeerisiä signaaleja, kun: Haluat yksinkertaisempaa äänestyspohjaista aistintaa, automaattista puhdistusta tai yhdistämistä toiminnan seuraamiseen.
public class RateLimitNotification : INotification
{
public int RetryAfterMs { get; init; }
}
await _mediator.Publish(new RateLimitNotification { RetryAfterMs = 5000 });
Paras: Eriytetty prosessinaikainen tapahtumakäsittely usean käsittelijän kanssa.
Käytä MediatR:ää, kun: Haluat useiden käsittelijöiden reagoivan samaan tapahtumaan synkronisesti.
Käytä efeerisiä signaaleja, kun: Haluat ympäristöntunnistusta ilman nimenomaista tilausta, itsepuhdistavaa historiaa tai integroitumista rajoitettuun suoritukseen.
var circuitBreaker = Policy
.Handle<HttpRequestException>()
.CircuitBreakerAsync(5, TimeSpan.FromSeconds(30));
Paras: Resilienssi yksittäisten puheluiden ympärillä automaattisella valtionhallinnalla.
Käytä Pollya, kun: Tarvitset puhelukohtaista häiriönsietokykyä automaattisilla puoli-avoimilla/suljetuilla siirtymillä.
Käytä efeerisiä signaaleja, kun: Haluat ympäristötietoisuuden monista toiminnoista, mukautetun piirilogiikan tai integroitumisen toiminnan seurantaan.
Yhdistä ne: Käytä Pollya työkehossasi, lähetä signaaleja virtapiireissä.
Lähestymistapa Itsepuhdistava Ympäristötietoisuus Eriyttäytyminen Custom Logic Integration |----------|:-------------:|:---------------:|:---------:|:------------:|:-----------:| OpenTelemetria Ulkoiset työkalut Reaktiiviset laajennukset Kompleksit Mediatr, käsikirja Polly Circuit Breaker Per call | Ekseeriset signaalit Sisäänrakennettu talo
Toimien toteutus 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; }
}
Työkehon sisällä:
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;
}
});
Signaalit ovat vain naruja. Käytä yksinkertaisia nimiä ("rate-limit") tai strukturoituja nimiä ("rate-limit:5000ms").
Mallisuodattimet käyttävät pallosemantiikkaa (*, ?) ja tukea pilkkuluetteloita ("error.*,timeout"). Täsmääminen on determinististä ja jakovaloa kautta StringPatternMatcher.
Hyvin yksityiskohtaista havainnointia varten voit lähettää signaaleja toiminnan jokaisessa vaiheessa. Kirjastossa on näyte SignalingHttpClient joka osoittaa tämän kaavan:
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...
});
Tämä lähettää signaaleja jokaisessa vaiheessa:
Signaali milloin?
|--------|------|
| stage.starting Ennen kuin pyyntö alkaa
| progress:0 Alkuvaiheen edistysmerkki
| stage.request HTTP:n pyyntö lähetetty
| stage.headers Vastausten otsikot vastaanotettiin
| stage.reading Alkaa lukea kehoa
| progress:XX Progress-prosentti (0-100) latauksen aikana
| stage.completed Lataus valmis
Voit sitten tiedustella näitä kaavalla, joka vastaa:
// 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
}
Toiminta voi myös poistaa omia signaalejaan. Tästä on hyötyä väliaikaisille valtioille:
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;
}
});
Signaalilähetyksen tavoin takaisinkytkennät voivat laukaista takaisinkutsuja:
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);
}
});
Erytropoietiini SignalRetractedEvent sisältää:
Signal - Signaalin nimi on peruttu.OperationId - Leikkaus, joka perui sen.Key - Operaation avain (jos sellainen on)Timestamp - Peruuttamisen yhteydessäWasPatternMatch - Totta, jos se vedetään takaisin RetractMatchingPattern - Käytetty kaava (jos kuvio täsmää)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;
}
});
Kaikki koordinaattorit antavat optimoidun signaalitiedustelun:
// 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");
From EphemeralOptions.cs:
Koordinaattorit voivat reagoida automaattisesti signaaleihin:
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)
});
Kun signaali tulee CancelOnSignals on havaittu, uudet kohteet ohitetaan (lasketaan epäonnistuneiksi).
Kun signaali tulee DeferOnSignals on havaittu, uudet kohteet odottavat, kunnes signaali hälvenee.
From 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;
}
}
Jokainen tämän palvelun ilmentymä peruuntuu automaattisesti, kun korkorajoihin isketään. Ei jaettua tilaa. Ei viestien kulkua. Luetaan vain hetkellistä ikkunaa.
Useita koordinaattoreita voi aistia toisensa jaetun kautta 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")
}
});
}
Mittarikirjastoa ei tarvita. Kysele vain hetkelliseltä ikkunalta.
From 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);
Virtapiirin katkaisimella ei ole omaa tilaa, se vain lukee ohikiitävää ikkunaa.
From Signaalit.cs:
Kun signaalit voivat aiheuttaa muita signaaleja, vaarana on ääretön silmukka. SignalConstraints estää tämän:
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);
}
}
};
Radan syy-seuraussuhde 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);
}
}
Lisääntymisketju seuraa polkua: order-placed → inventory-reserved → ...
Jos inventory-reserved yritti päästää irti order-placed, se olisi tukossa (sykli havaittu).
From Signaalit.cs:
Signaalit, joiden on oltava näkyviä koordinaattoreille:
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; }
}
Käyttö:
// 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();
}
From StringPatternMatcher.cs:
Glob-tyylinen vastaavuus signaalin suodatukseen:
// 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
}
}
Pidä signaalit yksinkertaisina ja johdonmukaisina:
// 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");
Synkroniset signaalien käsittelijät (OnSignal) Suorita toiminnon kierteellä – pidä ne nopeina. I/O-työn tekemiseen tuuletin lähettää signaalin async-polulle, jossa SignalDispatcher (mallien täsmäytys, deterministinen järjestys) tai 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
});
Mallit tukevat *, ?, ja pilkkulistat ("error.*,timeout"). Kaikki vastaavat käsittelijät toimivat rekisteröintijärjestyksessä tausta-avaimella varustetulla koordinaattorilla; päästöt pysyvät synkronisina.
Erillisessä asynkkikäsittelyssä:
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));
Täydellinen esimerkki, jossa yhdistetään asynkkisignaalien käsittely ja telemetriaintegraatio:
Lähde: TelemetriaSignalHandler.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();
}
Lähetä se koordinaattorillesi:
await using var telemetryHandler = new TelemetrySignalHandler(telemetryClient);
await using var coordinator = new EphemeralWorkCoordinator<Request>(
ProcessAsync,
new EphemeralOptions
{
OnSignal = signal => telemetryHandler.OnSignal(signal)
});
Käsittelijä:
OnSignal palaa hetiSignaalit ovat voimakkaita, koska ne ovat hetkellinen:
Omaisuus Hyödyt Hyödyt Hyödyt Hyödyt |----------|---------| | Rajattu Ei voi kasvaa ilman rajoja - vanhat signaalit vanhenevat | Itsepuhdistava Siivouskoodia ei tarvita | Tuotannosta irrotettu Lähettäjät eivät tiedä kuulijoista | Havaittavissa Mikä tahansa koodi voi aistia nykytilan | Yksityinen Ei käyttäjätietoja - vain signaalinimet | Nopea O(1) Havaitseminen oikosulkujen avulla
Effeeraalinen ikkuna on jo olemassa vianetsintää varten. Signaalit antavat sille vain semanttisen merkityksen.
Signaalit muuttavat eristetyt teloitusatomit AistinverkkoKukin koordinaattori voi:
CancelOnSignals sekä DeferOnSignalsEi viestimeklaria, ei yhteistä osavaltiota, ei koordinaatioprotokollaa, vain luonnollisesti hajoavia metatietoja.
Atomit eivät puhu toisilleen suoraan - ne vain jättävät jälkiä hetkelliseen ikkunaan, jota muut voivat tarkkailla. Se on stigmergiaa async-järjestelmille.
Tuli, signaali, aisti, unohda.
© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.