Back to "Efeeriset signaalit - Atomien muuttaminen sensoriverkoksi"

This is a viewer only at the moment see the article on how this works.

To update the preview hit Ctrl-Alt-R (or ⌘-Alt-R on Mac) or Enter to refresh. The Save icon lets you save the markdown file to disk

This is a preview from the server running through my markdig pipeline

Architecture ASP.NET Async Systems Design

Efeeriset signaalit - Atomien muuttaminen sensoriverkoksi

Friday, 12 December 2025

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.

NUGET!!!

Tämä on nyt enimmäkseen lucid.efemerals Nuget-paketti myös yli 20 enimmäkseen lucid.efemeral-kuviota ja 'atoms'.

NuGet Lisenssi

Lähdetiedostot

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

Ongelma: Yksinäiset atomit

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:

  • "Tämä API vain nopeusrajoitettu"
  • "Tämä käyttäjä on epäonnistunut kolme kertaa"
  • "Sisarusoperaatio odottaa yhä porttia"

Koordinaattorit voivat sitten muuttaa käytöstään ikkunassa näkyvien signaalien perusteella. Se on ympäristön tiedostamista ilman riippuvuuksia.


Ratkaisu: Toimintaa koskevia viestejä

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.


Efeeriset signaalilait

Laki 1 - Signaalit eivät koskaan implisiittisesti käynnistä teloitusta

Signaali itsessään ei aiheuta teloitusta. Se tallentaa vain faktan hetkellisessä ikkunassa. Mikään ei kulje, koska signaali lähetettiin.

Laki 2 - OnSignal on eksplisiittinen, paikallinen ja valinnainen

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.

Laki 3 - Atomin signaalit ovat paikallisia

Signaali kiinnittyy vain sen säteilevään atomiin/toimintoon. Mikään signaali ei koskaan mutaatioita tai anota toista atomia. Ei ole yhteistä kirjoitettavaa bussia.

Laki 4 - Viestit ovat vain faktoja

Signaalin lähettäminen lisää tosiasian atomien historiaan. Signaaleja ei koskaan päivitetä tai ylikirjoiteta. Peruutukset poistavat vain lähettimien omat signaalit.

Laki 5 - Signaalipinta on sidottu ja ajastettu

Signaaleja on vain koordinaattorin hetkellisessä ikkunassa. Ne vanhenevat automaattisesti, kun ikkuna vanhenee. Mitään ei itsepintaisesti tehdä, jos sinnikkyyttä ei nimenomaan rakenneta.

Laki 6 - Tarkkailijat lukevat, he eivät kirjoita

Kun koordinaattori tarkistaa Kööpenhamina-avaimen signaalit, se skannataan:

atomit ikkunassaan

ja paikalliset signaalit noille atomeille Tarkkailijat eivät koskaan muokkaa atomien signaalitilaa.

Laki 7 - Tapahtumamainen käytös on kerros, ei alkukantainen

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.

Laki 8 - Ei ristiatomimutaatiota missään olosuhteissa

Mikään atomi ei saa muuttaa toisen atomin kuvaa, tilaa, signaaleja tai metadataa. Yhteensovittaminen tapahtuu seuraavasti:

  • signaaleja

  • aistiminen

  • ikkunat

  • politiikka

Ei kirjoitusten kautta.

Laki 9 - Maailmanlaajuisia näkemyksiä johdetaan, ei koskaan mutageenia

SignalSinkin tai parven erittelemällä pinnalla voi olla yhdistetty kuva signaaleista – Mutta se on aina vain luettavaa, ei koskaan arvovaltaista, ei koskaan kirjoitettavaa.

Laki 10 - Käsittelijöiden poistaminen tuottaa saman ydinjärjestelmän

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.

Paras analogia: jalanjäljet, ei ohjeita

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.

Mitä tämä poistaa

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

Nopea vertailu

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

Koodien ero

Tapahtumia koskeva lähestymistapa (klassinen):

public event Action RateLimited;

try
{
    await CallApiAsync();
}
catch (RateLimitException)
{
    RateLimited?.Invoke(); // Makes someone else act right now
}

Ongelmia:

  • Kuka reagoi?
  • Missä järjestyksessä?
  • Entä jos kaksi käsittelijää on ristiriidassa keskenään?
  • Entä jos käsittelijä heittää?
  • Entä jos nyt on väärä aika?
  • Entä jos soittaja ei halua tällaista käytöstä?

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:

  • Ei suoraa kytkentää
  • Ei yllätyssoittoja
  • Ei teloitushyppyjä
  • Ei takaisinkutsua helvetisti
  • Ei ristikkäistä outoutta
  • Reagoi vain, kun kysy katsoa

"Aha!" -linja

Tapahtumasiirron ohjaus, viestien siirtokonteksti.

Siinä on koko henkinen malli yhdessä lauseessa.

Millaisia viestejä todellisuudessa on?

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.


Miten tämä on verrattavissa muihin lähestymistapoihin

Application Insights / OpenTelemetria

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.

Reaktiiviset laajennukset (Rx)

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.

Mediatr-ilmoitukset

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.

Polly Circuit Breaker

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ä.

Vertailutaulukko

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


Signaalien nostattaminen ja vetäminen takaisin

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.

Hienopiirteinen signaaliesimerkki: HTTP-puhelut

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
}

Vetävät signaalit

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

Peruutustapahtumat

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 RetractMatching
  • Pattern - Käytetty kaava (jos kuvio täsmää)

Reaalimaailman esimerkki: Rate Limit Recovery

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

Sensuroivia viestejä

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

Signaalin reagointi

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.


Real-World Example: Adaptive Rate Limitting

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.


Real-World Example: Cross-Coordinator Awareness

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

Real-World Example: Health Monitoring

[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.


Real-World Example: Signaalipohjainen piirinmurto

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.


Signaalirajoitteet: Äärettömien silmukoiden estäminen

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

Signaalin välittäminen

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


SignalSink: Global Signal Space

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

StringPatternMatcherin malli

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

Signaalin nimeäminen konventeissa

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

Asynkin signaalinkäsittely

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.

SignalDispatcherin uloskirjautuminen

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.

AsyncSignalProcessor

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

TelemetriaSignalHandler -esimerkki

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ä:

  • Älä koskaan estä leikkauslanka - OnSignal palaa heti
  • Luokittelee signaaleja eri telemetriatyypeille etuliitteellä
  • Paljastaa metriikan (kyseessä, käsitelty, pudotettu) terveyden seurantaan
  • Rajaa muistia - laskee vanhimpia signaaleja, jos jono täyttyy

Miksi tämä toimii

Signaalit 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.


Päätelmät

Signaalit muuttavat eristetyt teloitusatomit AistinverkkoKukin koordinaattori voi:

  • Emit signaaleja siitä, mitä se koki
  • Järkevä signaaleja omasta historiasta tai yhteisestä nielusta
  • Reaktio automaattisesti kautta CancelOnSignals sekä DeferOnSignals

Ei 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.


Linkkejä

logo

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