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.
Questo è ora nel pacchetto per lo piùlucid.efemerals Nuget anche più di 20 modelli per lo più lucidi.efemeri e 'atomi'.
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 |
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:
I coordinatori possono quindi cambiare il loro comportamento in base ai segnali visibili nella loro finestra. Questa è consapevolezza ambientale senza dipendenze.
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.
Un segnale, da solo, non provoca l'esecuzione. Registra solo un fatto nella finestra effimera. Niente funziona perche' e' stato emesso un segnale.
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.
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.
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.
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.
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.
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.
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.
Un lavello o una superficie aggregata a sciame può mostrare una visione combinata dei segnali ma è sempre di sola lettura, mai autorevole, mai scrivibile.
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.
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.
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 |
| 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 |
Approccio agli eventi (classico):
public event Action RateLimited;
try
{
await CallApiAsync();
}
catch (RateLimitException)
{
RateLimited?.Invoke(); // Makes someone else act right now
}
Problemi:
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:
Controllo del trasferimento degli eventi, contesto di trasferimento dei segnali.
E' l'intero modello mentale in una frase.
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.
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.
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.
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.
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.
| Avvicinamento | Autopulente | Sensing ambiente | Discoppiato | Logica personalizzata | Integrazione |
|---|---|---|---|---|---|
| Telemetria aperta | • • • | • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • | |||
| Estensioni reattive | ' ' ' ' ' ' ' ' ' ' ' ' ' ' ' ' ' complesso ' ' | ||||
| MediatR | ' ' ' ' ' ' ' ' ' ' ' ' | Manuale ' ' | |||
| Interruttore di interruttore a sfera | ' ' ' ' ' ' ' ' ' ' ' ' ' ' ' per chiamata ' ' | ||||
| Segnali effimeri • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • • |
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.
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
}
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;
}
});
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 ritrattoOperationId - L'operazione che l'ha ritrattata.Key - La chiave dell'operazione (se presente)Timestamp - Quando si è verificata la ritrazioneWasPatternMatch - Vero se ritratto via RetractMatchingPattern - Il modello utilizzato (se corrisponde al modello)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;
}
});
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");
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.
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.
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);
}
}
[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.
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.
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);
}
}
};
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).
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();
}
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
}
}
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");
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.
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.
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));
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:
OnSignal ritorna immediatamenteI 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.
I segnali trasformano atomi di esecuzione isolati in un rete di rilevamento. Ogni coordinatore può:
CancelOnSignals e DeferOnSignalsNessun 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.
© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.