Segnali effimeri - Trasformare gli atomi in una rete di sensori (Italiano (Italian))

Segnali effimeri - Trasformare gli atomi in una rete di sensori

Friday, 12 December 2025

//

24 minute read

Un minuscolo primitivo che trasforma il lavoro concomitante in un sistema adattivo coordinato.

"Il modello dei segnali effimeri"

Dentro Parte 1 abbiamo costruito l'esecuzione effimera - limitato, privato, auto-pulire flussi di lavoro asincroni. Parte 2 L'abbiamo trasformata in una libreria riutilizzabile con coordinatori, condotte chiavi in mano e integrazione DI.

Questo articolo aggiunge una piccola caratteristica che cambia tutto: segnali.

NUGET!!!

Questo è ora nel pacchetto per lo piùlucid.efemerals Nuget anche più di 20 modelli per lo più lucidi.efemeri e 'atomi'.

NuGetCity name (optional, probably does not need a translation) Licenza

File sorgente

Il codice sorgente completo è nel per lo piùlucid.atoms Repository GitHub

L'infrastruttura del segnale risiede in:

File Scopo
Segnali.cs SignalEvent, SignalRetractedEvent, SignalPropagation, SignalConstraints, SignalSink, ISignalEmitter, AsyncSignalProcessor
Operazione effimera.cs Emissione del segnale e ritrazione dalle operazioni
Opzioni effimere.cs Configurazione reattiva al segnale (CancelOnSignals, DeferOnSignals, OnSignal, OnSignalRetracted)
StringPatternMatcher.cs Modello in stile Glob per il filtraggio del segnale
SignalDispatcher.cs Instradamento del segnale asincrono con corrispondenza del motivo (supporta *, ?, liste di virgola, ordine deterministico)
Esempi/SignalingHttpClient.cs Emissione del segnale a grana fine del campione per le chiamate HTTP
Esempi/AdaptiveTranslationService.cs Limitazione della velocità di adattamento con rinvio basato sul segnale
Esempi/SignalBasedCircuitBreaker.cs Finestra del segnale effimero di lettura dell'interruttore del circuito
Esempi/TelemetriaSignalHandler.cs Elaborazione del segnale asincrono con integrazione della telemetria

Il problema: atomi isolati

I nostri coordinatori effimeri sono bravi a lavorare, ma sono isolati. Ogni coordinatore conosce le proprie operazioni, ma non ha alcuna consapevolezza di ciò che sta accadendo altrove nel sistema.

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

Potremmo creare dipendenze esplicite, ma questo crea accoppiamento. sensibilizzazione ambientale - coordinatori che possono percepire il loro ambiente senza essere direttamente collegati.

I segnali permettono agli atomi di esecuzione di lasciare tracce nella loro finestra effimera. Quelle tracce si comportano come fatti di breve durata:

  • "Questa API limitata al tasso"
  • "Questo utente ha fallito 3 volte"
  • "Un'operazione tra fratelli sta ancora aspettando il gateway"

I coordinatori possono quindi cambiare il loro comportamento in base ai segnali visibili nella loro finestra. Questa è consapevolezza ambientale senza dipendenze.


La soluzione: Segnali sulle operazioni

Da Segnali.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);
}

Un'operazione può sollevare segnali durante l'esecuzione. Questi segnali vivono nella finestra effimera accanto all'operazione. Quando l'operazione invecchia, i segnali vanno con esso.

Niente broker di messaggi, niente infrastruttura separata, solo legami con le operazioni.

Poiché i segnali vivono all'interno della finestra effimera, ereditano le sue garanzie: dimensioni limitate, invecchiamento automatico e ciclo di vita zero.


Le leggi del segnale effimero

Legge 1 - I segnali non attivano mai in modo implicito l'esecuzione

Un segnale, da solo, non provoca l'esecuzione. Registra solo un fatto nella finestra effimera. Niente funziona perche' e' stato emesso un segnale.

Legge 2 - OnSignal è esplicito, locale e facoltativo

Se un coordinatore definisce un gestore OnSignal, viene eseguito in modo sincrono quando viene emesso un segnale - ma solo perché il coordinatore ha scelto di collegarlo. Gli emitter non sanno e non si preoccupano. La rimozione di tutti i gestori lascia invariato il comportamento del nucleo.

Legge 3 - I segnali sono locali all'atomo

Un segnale è collegato solo all'atomo/operazione che l'ha emesso. Nessun segnale muta o annota un altro atomo. Non c'è nessun bus scrivibile condiviso.

Legge 4 - I segnali sono solo fatti dell'appendice

Emettendo un segnale si aggiunge un fatto alla storia degli atomi. I segnali non vengono mai aggiornati o sovrascritti. Le ritrazioni rimuovono solo i segnali propri degli emettitori.

Legge 5 - La superficie del segnale è delimitata e limitata nel tempo

I segnali esistono solo all'interno della finestra effimera dei coordinatori. Scadono automaticamente quando la finestra invecchia. Niente è persistente a meno che non si costruisca esplicitamente persistenza.

Legge 6 - Gli osservatori leggono, non scrivono

Quando un coordinatore controlla i segnali per la chiave KFBG, esegue la scansione:

gli atomi nella sua finestra

e i segnali locali a quegli atomi Gli osservatori non modificano mai uno stato di segnale degli atomi.

Legge 7 - Comportamento simile a un evento è un livello, non un primitivo

Se vuoi che i segnali guidino flussi di lavoro asincroni, devi usare:

Dispatcher del segnale

AsyncSignalProcessor

o altri adattatori.

Si tratta di livelli opzionali costruiti in cima ai segnali, non parte della loro semantica.

Legge 8 - No Cross-Atom Mutation in Qualsiasi circostanza

Nessun atomo può alterare l'istantanea, lo stato, i segnali o i metadati di un altro atomo. Il coordinamento avviene attraverso:

  • segnali

  • rilevamento

  • finestre

  • politiche

Non attraverso gli scritti.

Legge 9 - Le viste globali sono derivate, mai mutabili

Un lavello o una superficie aggregata a sciame può mostrare una visione combinata dei segnali ma è sempre di sola lettura, mai autorevole, mai scrivibile.

Legge 10 - Rimozione manigliere Rende lo stesso sistema centrale

Se tutti i gestori (OnSignal, dispacciatori, processori) sono distaccati, il sistema rimane completamente corretto e prevedibile. I segnali hanno ancora un significato perché sono fatti, non trigger.

La migliore analogia: impronte, non istruzioni

Manifestazioni sono come una telefonata:

"Ti sto chiamando, rispondi e reagisci."

Segnali sono come impronte sulla neve:

"Ho lasciato le impronte. Se volete sapere dove sono andato - guardate. Se non vi interessa - ignoratelo."

Questo è il motivo per cui i segnali non si rompono mai, non bloccano mai, e non interagiscono mai con il flusso di controllo a meno che non si scegli per interrogarli.

Cosa rimuove

La distinzione è importante perché i segnali eliminano:

Problema Eventi hanno esso Segnali evitarlo
Accoppiamento Editore → Sottoscrivatori Nessuno
Dipendenze di tempo Richiesta reazione immediata Ambiente, sondaggio quando pronto
Loop di ritorno I gestori possono attivare i gestori Nessun loop a meno che non sia esplicitamente richiesto
Semantica d'ordine Materie d'ordine Irrilevante
Garantie di consegna Devono essere consegnate/gestite Nessuna consegna, solo esistenza
Propagazione degli errori Errori del gestore si propagano Isolato
Rischi di resistenza Comune Impossibile
Fallimenti cascati Un gestore fallisce, interruzioni della catena Nessuna catena da rompere

Confronto rapido

Caratteristica Eventi Segnali
Accoppiamento Forte (editore → abbonati) Nessuna
Timing Immediate Ambient
Consegna Garantito / tentato Nessuna consegna, solo esistenza
Reazione Richiesto Optional
Carico di pagamento dati Spesso pesante Meccanismi di piccole stringhe di metadati
Archiviazione Nessuna Finestra scorrevole in stile LRU
Tempo di vita Instantaneo Decavi automaticamente
Modalità di guasto Molti Quasi nessuno

La differenza di codice

Approccio agli eventi (classico):

public event Action RateLimited;

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

Problemi:

  • Chi reagisce?
  • In che ordine?
  • E se due gestori fossero in conflitto?
  • E se lanciasse un supervisore?
  • E se fosse il momento sbagliato?
  • E se il chiamante non volesse questo comportamento?

Avvicinamento del segnale (efemero):

try
{
    await CallApiAsync(ct);
}
catch (RateLimitException ex)
{
    op.Signal($"rate-limit:{ex.RetryAfterMs}ms");
    throw;
}

Non succede niente qui. Più tardi, in un posto completamente diverso:

if (translator.HasSignal("rate-limit"))
{
    await Task.Delay(1000);   // Act when *we* choose
}

Note:

  • Nessun accoppiamento diretto
  • Nessuna chiamata a sorpresa
  • Nessun salto di esecuzione
  • Nessun inferno di richiamo
  • Nessuna stranezza tra i fili
  • Reagisce solo quando noi chiedi per guardare

La linea "Aha!"

Controllo del trasferimento degli eventi, contesto di trasferimento dei segnali.

E' l'intero modello mentale in una frase.

Quali segnali sono realmente

A seconda del tuo background:

Udienza Definizione
Pensatori di sistemi Un substrato stigmergico per la coordinazione indiretta
Ingegneri metadati leggeri collegati alle operazioni in una finestra scorrevole delimitata
Nerds PL/Concurrency Una lavagna temporale implicita accoppiata a semantica della valuta delimitata
Utenti del quadro Una superficie di stato autopulente in-processo è possibile interrogare in qualsiasi momento

I segnali sono tracce lasciate in una superficie di memoria condivisa e delimitata. Chiunque può guardarli. Nessuno è obbligato a reagire. stigmergic modello di coordinamento - lo stesso che le formiche usano, lo stesso sistema di lavagna utilizzato nei primi AI, e lo stesso che le moderne reti di gossip CRDT suggeriscono.


Come questo si confronta con altri approcci

Approfondimenti dell'applicazione / OpenTelemetry

using var activity = source.StartActivity("ProcessOrder");
activity?.SetTag("order.id", orderId);
activity?.SetTag("rate.limited", true);

Meglio per: Tracciamento distribuito tra i servizi, storage telemetrico a lungo termine, ID di correlazione.

Usare la telemetria quando: È necessario tracciare le richieste attraverso più servizi, memorizzare metriche per l'analisi, o integrare con strumenti di monitoraggio.

Usare segnali effimeri quando: Hai bisogno di consapevolezza ambientale in-processo, coordinamento reattivo, o non vuole infrastrutture di telemetria.

Estensioni attive (Rx)

var rateLimits = Observable.FromEventPattern<RateLimitEventArgs>(
    h => api.RateLimitHit += h,
    h => api.RateLimitHit -= h);

rateLimits
    .Throttle(TimeSpan.FromSeconds(1))
    .Subscribe(e => HandleRateLimit(e));

Meglio per: Elaborazione complessa di eventi, operazioni basate sul tempo, combinando più flussi di eventi.

Usare Rx quando: Hai bisogno di interrogazioni temporali complesse (windowing, debouncing, combinando flussi).

Usare segnali effimeri quando: Si desidera più semplice rilevamento basato sul polling, pulizia automatica, o l'integrazione con il monitoraggio delle operazioni.

Notifiche MediatR

public class RateLimitNotification : INotification
{
    public int RetryAfterMs { get; init; }
}

await _mediator.Publish(new RateLimitNotification { RetryAfterMs = 5000 });

Meglio per: Gestione di eventi in-processo disaccoppiato con più gestori.

Usa MediatR quando: Vuoi che più gestori reagiscano allo stesso evento in modo sincrono.

Usare segnali effimeri quando: Desiderate il rilevamento ambientale senza l'abbonamento esplicito, la cronologia autopulente, o l'integrazione con l'esecuzione limitata.

Interruttore di circuito Polly

var circuitBreaker = Policy
    .Handle<HttpRequestException>()
    .CircuitBreakerAsync(5, TimeSpan.FromSeconds(30));

Meglio per: Resilienza intorno alle chiamate individuali con gestione automatica dello stato.

Usa Polly quando: Hai bisogno di resilienza per chiamata con transizioni automatiche semi-aperte/chiuse.

Usare segnali effimeri quando: Desiderate consapevolezza ambientale attraverso molte operazioni, logica del circuito personalizzato, o integrazione con il monitoraggio delle operazioni.

Combinarli: Utilizzare Polly all'interno del corpo di lavoro, emettere segnali quando i circuiti di viaggio.

Tabella di confronto

Avvicinamento Autopulente Sensing ambiente Discoppiato Logica personalizzata Integrazione
Telemetria aperta • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • •
Estensioni reattive ' ' ' ' ' ' ' ' ' ' ' ' ' ' ' ' ' complesso ' '
MediatR ' ' ' ' ' ' ' ' ' ' ' ' Manuale ' '
Interruttore di interruttore a sfera ' ' ' ' ' ' ' ' ' ' ' ' ' ' ' per chiamata ' '
Segnali effimeri • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • •

Alza e ritira i segnali

Attuazione delle operazioni 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; }
}

All'interno del corpo di lavoro:

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

I segnali sono solo stringhe. Usa nomi semplici ("rate-limit") o nomi strutturati ("rate-limit:5000ms"). I filtri modello usano la semantica glob (*, ?) e liste di virgole di supporto ("error.*,timeout"). L'abbinamento è deterministico e la luce di allocazione via StringPatternMatcher.

Esempio di segnale granato fine: chiamate HTTP

Per un'osservazione molto dettagliata, è possibile emettere segnali in ogni fase di un'operazione. La libreria include un campione SegnalazioneHttpClient che dimostra questo modello:

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

Questo emette segnali in ogni fase:

Segnale Quando
stage.starting Prima dell'inizio della richiesta
progress:0 Marchio iniziale di avanzamento
stage.request Richiesta HTTP inviata
stage.headers Intestazioni di risposta ricevute
stage.reading Avvio della lettura del corpo
progress:XX Percentuale di avanzamento (0-100) durante il download
stage.completed Download finito

Puoi quindi interrogarli con il modello corrispondente:

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

Ritirare i segnali

Le operazioni possono anche rimuovere i propri segnali. Questo è utile per gli stati temporanei:

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

Eventi di ritrazione

Proprio come l'emissione del segnale, le ritrazioni possono innescare i richiami:

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

La SignalRetractedEvent comprende:

  • Signal - Il nome del segnale ritratto
  • OperationId - L'operazione che l'ha ritrattata.
  • Key - La chiave dell'operazione (se presente)
  • Timestamp - Quando si è verificata la ritrazione
  • WasPatternMatch - Vero se ritratto via RetractMatching
  • Pattern - Il modello utilizzato (se corrisponde al modello)

Esempio di Real-World: 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;
    }
});

Segnali sensoriali

Tutti i coordinatori forniscono un'interrogazione ottimizzata del segnale:

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

Elaborazione del segnale-reattivo

Da Opzioni effimere.cs:

I coordinatori possono reagire automaticamente ai segnali:

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

Quando un segnale entra CancelOnSignals viene rilevato, i nuovi elementi vengono saltati (contati come falliti). Quando un segnale entra DeferOnSignals viene rilevato, i nuovi elementi attendono che il segnale si ripulisca.


Esempio di Real-World: Limitazione dei tassi di adattamento

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

Ogni istanza di questo servizio si disattiva automaticamente quando i limiti di velocità vengono colpiti. Nessun stato condiviso. Nessun messaggio che passa. Solo leggendo la finestra effimera.


Esempio di Real-World: Consapevolezza cross-coordinatore

Molti coordinatori possono percepirsi l'un l'altro attraverso una condivisione 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);
    }
}

Esempio di Real-World: Monitoraggio della salute

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

Non servono librerie metriche, basta interrogare la finestra effimera.


Esempio di Real-World: Breaking del circuito basato sul segnale

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

L'interruttore non ha un proprio stato - legge solo la finestra effimera.


Vincoli del segnale: Prevenire i loops infiniti

Da Segnali.cs:

Quando i segnali possono causare altri segnali, si rischia loop infiniti. SignalConstraints impedisce che:

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

Propagazione del segnale

Track causality with 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);
    }
}

La catena di propagazione traccia il percorso: order-placed → inventory-reserved → ...

Se inventory-reserved cercato di emettere order-placed, sarebbe bloccato (ciclo rilevato).


Il lavello del segnale: spazio globale del segnale

Da Segnali.cs:

Per i segnali che devono essere visibili tra i coordinatori:

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

Uso:

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

Abbinamento modello con StringPatternMatcher

Da StringPatternMatcher.cs:

Corrispondenza in stile Glob per il filtraggio dei segnali:

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

Convenzioni di denominazione del segnale

Mantenere i segnali semplici e coerenti:

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

Gestione del segnale asincrono

Gestori di segnali sincroni (OnSignal) eseguire sul thread dell'operazione li mantiene veloci. Per fare il lavoro I/O-bound, i segnali della ventola in un percorso asincrono con SignalDispatcher (corrispondente al modello, ordine deterministico) o AsyncSignalProcessor.

Fan-out SignalDispatcher

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

Supporto dei modelli *, ?, e liste di virgola ("error.*,timeout") Tutti i gestori di corrispondenza funzionano in ordine di registrazione su un coordinatore con tasto di sfondo; l'emissione rimane sincrona.

AsyncSignalProcessor

Per l'elaborazione asincrona standalone:

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

Esempio di telemetriaSignalHandler

Un esempio completo che combina l'elaborazione del segnale asincrono con l'integrazione della telemetria:

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

Collegalo al tuo coordinatore:

await using var telemetryHandler = new TelemetrySignalHandler(telemetryClient);

await using var coordinator = new EphemeralWorkCoordinator<Request>(
    ProcessAsync,
    new EphemeralOptions
    {
        OnSignal = signal => telemetryHandler.OnSignal(signal)
    });

Il responsabile:

  • Mai blocchi il filetto di funzionamento - OnSignal ritorna immediatamente
  • Categorizza i segnali per prefisso per diversi tipi di telemetria
  • Espone le metriche (liquidazione, trasformazione, riduzione) per il monitoraggio sanitario
  • Memoria di bordi - lascia cadere i segnali più vecchi se la coda si riempie

Perché funziona

I segnali sono potenti perché sono effimero:

Proprietà Benefit
Confinato Non può crescere slegato - vecchi segnali età fuori
Autopulente Nessun codice di pulizia necessario
disaccoppiata Emitters don't know about listeners
Osservabile Qualsiasi codice può percepire lo stato corrente
Privato Nessun dato utente - solo nomi di segnale
Veloce Rilevamento O(1) con cortocircuito

La finestra effimera è già lì per il debug. I segnali danno solo significato semantico.


Conclusione

I segnali trasformano atomi di esecuzione isolati in un rete di rilevamento. Ogni coordinatore può:

  • EmitCity name (optional, probably does not need a translation) segnali su ciò che ha sperimentato
  • Senso segnali dalla propria storia o da un lavandino condiviso
  • Reagisci automaticamente tramite CancelOnSignals e DeferOnSignals

Nessun broker di messaggi, nessun stato condiviso, nessun protocollo di coordinamento, solo operazioni con metadati che naturalmente decadranno.

Gli atomi non si parlano direttamente tra loro - lasciano solo tracce nella finestra effimera che altri possono osservare. E' stigmergia per i sistemi asincroni.

Fuoco... segnale... senso... scordatelo.


Collegamenti

Finding related posts...
logo

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