Back to "Efemerala signaler - att förvandla atomer till ett senserande nätverk"

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

Efemerala signaler - att förvandla atomer till ett senserande nätverk

Friday, 12 December 2025

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.

NUGET!!!

Detta är nu också i det mest lucid.ephemerals Nuget-paketet Mer än 20 oftast lucid.ephemerala mönster och 'atomer'.

Hämta Licens

Källfiler

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

Problemet: isolerade atomer

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:

  • "Detta API bara hastighetsbegränsad"
  • "Den här användaren har misslyckats 3 gånger"
  • "En syskonoperation väntar fortfarande på porten"

Koordinatorerna kan då ändra sitt beteende baserat på de signaler som syns i deras fönster. Det är omgivningsmedvetenhet utan beroenden.


Lösningen: Signaler om verksamheten

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.


De efemära signallagarna

Lag 1 - Signaler aldrig oförklarligt Trigger Execution

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.

Lag 2 - OnSignal är explicit, lokal och frivillig

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.

Lag 3 - Signaler är lokala för atomen

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.

Lag 4 - Signaler är bara bihang fakta

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.

Lag 5 - Signalytan är begränsad och tidsbegränsad

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.

Lag 6 - Observatörer läser, de skriver inte

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.

Lag 7 - Händelseliknande beteende är ett lager, inte en primitiv

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.

Lag 8 - Ingen tväratom Mutation under några omständigheter

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.

Lag 9 - Globala uppfattningar är härledda, aldrig muterbara

En signalSänka eller svärmaggregaterad yta kan visa en kombinerad vy av signaler — Men den är alltid skrivbar, aldrig auktoritativ, aldrig skrivbar.

Lag 10 - Ta bort Handledare ger samma kärnsystem

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.

Den bästa analogin: Fotavtryck, inte instruktioner

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.

Vad detta tar bort

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

Snabb jämförelse

. 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

Skillnaden i kod

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:

  • Vem reagerar?
  • I vilken ordning?
  • Tänk om två hanterare bråkar?
  • Tänk om en handläggare kastar?
  • Tänk om det är fel tidpunkt?
  • Tänk om den som ringer inte vill ha det här beteendet?

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:

  • Ingen direkt koppling
  • Inga överraskningssamtal
  • Ingen avrättning hoppar
  • Inget kalla tillbaka helvete
  • Ingen korstrådsskumhet
  • Reaktioner bara när vi fråga för att titta

Linjen "Aha!"

Överföringskontroll av händelser, signalöverföringskontext.

Det är hela mentalmodellen i en mening.

Vad signaler egentligen är

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


Hur detta kan jämföras med andra metoder

Application Insikter / OpenTelemetri

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.

Reaktiva förlängningar (Rx)

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.

MediatR-anmälningar

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.

Polly Circuit Breaker

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.

Jämförelsetabell

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.


Upphöjande och upplyftande signaler

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.

Exempel på finjusterad signal: HTTP-anrop

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
}

Upprullningssignaler

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

Återtagningshändelser

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 signalnamnet
  • OperationId - Operationen som drog tillbaka den.
  • Key - Operationens nyckel (om sådan finns)
  • Timestamp - När reträtt inträffade
  • WasPatternMatch - Sant om de dras in via RetractMatching
  • Pattern - Mönstret som används (om mönstret matchar)

Real-World-exempel: Hastighetsbegränsning Återvinning

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

Senserande signaler

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

Signalreaktiv bearbetning

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.


Real-World Exempel: Adaptive Rate Limiting

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.


Real-World Exempel: Cross-Coordinator Medvetenhet

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

Real-World Exempel: Hälsoövervakning

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


Real-World-exempel: Signalbaserad kretsbrytning

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.


Signalbegränsningar: Förhindra oändliga loopar

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

Signalförökning

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


Signalsänkan: Globalt signalutrymme

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

Mönster som matchar med strängmönsterMatcher

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

Signalnameringskonventioner

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

Asynkroniseringssignalhantering

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.

Signalutblåsning från Dispatcher

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.

AsyncSignalProcessorName

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

Exempel på telemetrisignalhandler

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:

  • Blockera aldrig insatstråden - OnSignal återvänder omedelbart
  • Kategoriserar signaler enligt prefix för olika telemetrityper
  • Mätvärden för exponering (kväde, bearbetade, tappade) för hälsoövervakning
  • Gränser minne - släpper de äldsta signalerna om kön fylls upp

Varför detta fungerar

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


Slutsatser

Signaler förvandlar isolerade avrättningsatomer till en Avkänning av nätverk. Varje samordnare kan:

  • Emit Ordförande signalerar om vad den upplevde
  • Känsla signaler från sin egen historia eller en delad diskho
  • Reagera automatiskt via CancelOnSignals och DeferOnSignals

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


Länkar

logo

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