Señales efímeras - Convertir átomos en una red de detección (Español (Spanish))

Señales efímeras - Convertir átomos en una red de detección

Friday, 12 December 2025

//

26 minute read

Un pequeño primitivo que convierte el trabajo concurrente en un sistema coordinado y adaptativo.

"El Patrón de Señales Efímeras"

In Parte 1 construimos la ejecución efímera - flujos de trabajo limitados, privados, auto-limpieza async. En Parte 2 lo convertimos en una biblioteca reutilizable con coordinadores, tuberías con llave, e integración DI.

Este artículo añade una pequeña característica que cambia todo: señales.

Esto está ahora en el paquete Nuget mayoritariamente lucid.efemerals también más de 20 patrones mayoritariamente lucid.efemerals y 'átomos'.

NuGet Licencia

Archivos de origen

El código fuente completo está en el mayormentelucid.atoms GitHub repositorio

La infraestructura de señales vive en:

Archivo # # Propósito

------ ---------
EfímeroOperación.cs Emisión de señales y retracción de las operaciones
EfímeroOpciones.cs Configuración de reacción a la señal (CancelOnSignals, DeferOnSignals, OnSignal, OnSignalRetracted)
StringPatternMatcher.cs Patrón estilo Glob emparejamiento para el filtrado de señales
SeñalDispatcher.cs Enrutamiento de la señal de Async con la coincidencia del patrón (soportes *, ?, listas de comas, orden determinista)
Ejemplos/SeñalizaciónHttpClient.cs Muestra de emisión de señal de grano fino para llamadas HTTP
Ejemplos/AdaptiveTranslationService.cs Limitación de la velocidad de adaptación con aplazamiento basado en la señal
Ejemplos/SeñalBasedCircuitBreaker.cs Interruptor de circuito leyendo la ventana de señal efímera
Ejemplos/TelemetríaSeñalHandler.cs Procesamiento de señales asinc con integración de telemetría

El problema: Átomos aislados

Nuestros coordinadores efímeros son excelentes para procesar el trabajo, pero están aislados. Cada coordinador sabe de sus propias operaciones, pero no tiene conciencia de lo que está sucediendo en otras partes del 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

Podríamos conectar dependencias explícitas, pero eso crea acoplamiento. Sensibilización ambiental - coordinadores que puedan percibir su entorno sin estar directamente conectados.

Las señales dejan que los átomos de ejecución dejen rastros en su ventana efímera.

  • "Esta API acaba de limitar la tasa"
  • "Este usuario ha fallado 3 veces"
  • "Una operación de hermano todavía está esperando en la puerta de entrada"

Los coordinadores pueden entonces cambiar su comportamiento basado en las señales visibles en su ventana. Eso es conciencia ambiental sin dependencias.


La solución: Señales sobre las operaciones

Desde Señales.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);
}

Una operación puede levantar señales durante la ejecución. Esas señales viven en la ventana efímera junto a la operación. Cuando la operación envejece, las señales van con ella.

Eso es, sin corredor de mensajes, sin infraestructuras separadas, sólo cadenas atadas a las operaciones.

Debido a que las señales viven dentro de la ventana efímera, heredan sus garantías: tamaño limitado, envejecimiento automático y sobrecarga del ciclo de vida cero.


Las leyes de la señal efímera

Ley 1 - Señales que nunca activan implícitamente la ejecución

Una señal, por sí misma, no causa la ejecución. Sólo registra un hecho en la ventana efímera. Nada funciona porque una señal fue emitida.

La Ley 2 - OnSignal es explícita, local y opcional

Si un coordinador define un controlador OnSignal, se ejecuta sincrónicamente cuando se emite una señal - pero sólo porque el coordinador optó por adjuntarla. Los emisores no saben ni se preocupan. La eliminación de todos los manipuladores deja sin cambios el comportamiento del núcleo.

Ley 3 - Las señales son locales al átomo

Una señal se une sólo al átomo/operación que la emitió. Ninguna señal muta o anota otro átomo. No hay bus compartido por escrito.

Ley 4 - Las señales son hechos que sólo se añaden

La emisión de una señal añade un hecho a la historia del átomo. Las señales nunca se actualizan o sobrescriben. Las retracciones eliminan sólo las señales propias del emisor.

Ley 5 - La superficie de la señal está limitada y limitada por el tiempo

Las señales sólo existen dentro de la ventana efímera del coordinador. Caducan automáticamente a medida que la ventana envejece. Nada persiste a menos que construyas explícitamente persistencia.

Ley 6 - Los observadores leen, no escriben

Cuando un coordinador “comprueba las señales de la clave K”, escanea:

los átomos en su ventana

y las señales locales a esos átomos Los observadores nunca modifican el estado de señal de un átomo.

Ley 7 - El comportamiento semejante a un evento es una capa, no una primitiva

Si desea que las señales impulsen flujos de trabajo asíncronos, debe utilizar:

SignalDispatcher

AsyncSignalProcessor

u otros adaptadores.

Estas son capas opcionales construidas encima de las señales, no parte de su semántica.

Ley 8 - No hay mutación de átomos cruzados bajo ninguna circunstancia

Ningún átomo puede alterar la instantánea, el estado, las señales o los metadatos de otro átomo. La coordinación se realiza mediante:

  • señales

  • detección

  • ventanas

  • políticas

No a través de las escrituras.

Ley 9 - Las opiniones globales son derivadas, nunca mutables

Una superficie de tinta o enjambre puede mostrar una visión combinada de las señales: pero siempre es de sólo lectura, nunca autorizado, nunca escrito.

Ley 10 - La eliminación de los manipuladores entrega el mismo sistema de núcleo

Si todos los manipuladores (OnSignal, despachadores, procesadores) están separados, el sistema sigue siendo totalmente correcto y previsible. Las señales todavía tienen significado porque son hechos, no desencadenantes.

La mejor analogía: Huellas, no instrucciones

Acontecimientos son como una llamada telefónica:

"Te estoy llamando ahora mismo.

Señales son como huellas en la nieve:

"Dejé huellas, si quieres saber a dónde fui, mira, si no te importa, ignóralas".

Esta es la razón por la que las señales nunca se rompen, nunca bloquean, y nunca interactúan con el flujo de control a menos que usted elegir para interrogarlos.

Lo que esto elimina

La distinción importa porque las señales eliminan:

Problemas # # Los eventos lo tienen # # Las señales lo evitan

|---------|:--------------:|:----------------:| Acoplamiento Editor → Suscriptores Ninguno Dependencias de tiempo Reacción inmediata requerida Ambiente, encuesta cuando esté listo bucles de retroalimentación Los manipuladores pueden desencadenar los manipuladores No hay bucles a menos que se pregunte explícitamente Ordenando semántica El orden importa Irrelevante Garantías de entrega Debe ser entregado / manejado Ninguna entrega, existencia justa Se propaga el error Se propaga el error Manipulador Aislado

Riesgos de reentrada # # comunes # imposibles

Fallas cascadas # # Un manejador falla, rompe cadenas # # No hay cadena que romper # # Un manejador falla, rompe cadenas # # No hay cadena que romper # # No hay cadena que romper

Comparación rápida

Acontecimientos Señales |---------|--------|---------| Acoplamiento Fuerte (editor → suscriptores) Ninguno

Sincronización # # Inmediata # # Ambient # # Ambient # # Ambient # # Ambient # # Ambient # # Ambient # # Ambient # # Ambient # # Ambient # # Ambient # Ambient # Ambient # Ambient # Ambient

Entrega # # Garantizada/intentada # # No entrega, sólo existencia

Reacción Requerida Opcional Carga útil de datos A menudo pesadas metadatos de cadenas diminutas

Almacenamiento # # Ninguno # Ventana deslizante de estilo LRU

La vida # # Instantánea # # Se decae automáticamente

Modos de fracaso # # Muchos # Casi ninguno

La diferencia en el código

Acontecimientos (clásicos):

public event Action RateLimited;

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

Problemas:

  • ¿Quién reacciona?
  • ¿En qué orden?
  • ¿Y si dos manejadores entran en conflicto?
  • ¿Qué pasa si un manejador lanza?
  • ¿Y si es el momento equivocado?
  • ¿Y si la persona que llama no quiere este comportamiento?

Aproximación de la señal (efemeral):

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

No pasa nada aquí, más tarde, en un lugar totalmente diferente.

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

Notas:

  • Sin acoplamiento directo
  • No hay llamadas sorpresa
  • Sin saltos de ejecución
  • Ningún infierno de devolución de llamadas
  • Sin rarezas en los hilos cruzados
  • Sólo reacciona cuando nosotros preguntar para mirar

La línea "Aha!"

Control de transferencia de eventos. Señales de contexto de transferencia.

Ese es todo el modelo mental en una frase.

¿Qué son realmente las señales?

Dependiendo de sus antecedentes:

audiencia definición |----------|------------| | Pensadores de sistemas Un sustrato estigérgico para la coordinación indirecta | Ingenieros Metadatos ligeros unidos a las operaciones en la ventana corredera limitada | PL/Nerds en moneda extranjera Una pizarra temporal implícita emparejada con semántica de concurrencia limitada | Usuarios del Marco Una superficie de estado en proceso, auto-limpieza que usted puede consultar en cualquier momento

Las señales son trazas que quedan en una superficie de memoria compartida y limitada. Cualquiera puede mirarlas. Nadie está obligado a reaccionar. estigérgico modelo de coordinación - el mismo uso de hormigas, el mismo sistema de pizarra utilizado en la IA temprana, y el mismo moderno CRDT chismes redes de insinuación.


Cómo se compara esto con otros enfoques

Insights de aplicaciones / OpenTelemetry

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

Lo mejor para: Rastreo distribuido a través de servicios, almacenamiento de telemetría a largo plazo, identificaciones de correlación.

Usar telemetría cuando: Necesita rastrear peticiones en múltiples servicios, almacenar métricas para el análisis o integrarse con herramientas de monitoreo.

Usar señales efímeras cuando: Usted necesita conciencia ambiental en el proceso, coordinación reactiva, o no quiere infraestructura de telemetría.

Extensiones reactivas (Rx)

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

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

Lo mejor para: Procesamiento de eventos complejos, operaciones basadas en el tiempo, combinando múltiples flujos de eventos.

Usar Rx cuando: Necesita consultas temporales complejas (rebobinación, desboneamiento, combinación de corrientes).

Usar señales efímeras cuando: Usted quiere una detección basada en encuestas más simple, limpieza automática, o integración con el seguimiento de la operación.

Notificaciones MediatR

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

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

Lo mejor para: Manipulación de eventos desacoplados en el proceso con múltiples manipuladores.

Usar MediatR cuando: Usted quiere que múltiples manejadores reaccionen al mismo evento sincrónicamente.

Usar señales efímeras cuando: Usted desea la detección ambiental sin suscripción explícita, historia de auto-limpieza, o integración con la ejecución limitada.

Interruptor de circuito Polly

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

Lo mejor para: Resiliencia en torno a llamadas individuales con gestión automática del estado.

Usar Polly cuando: Necesita resiliencia por llamada con transiciones automáticas semiabiertas/cerradas.

Usar señales efímeras cuando: Usted quiere conciencia ambiental a través de muchas operaciones, lógica de circuito personalizado, o integración con el seguimiento de la operación.

Combínalos.: Utilice Polly dentro de su cuerpo de trabajo, emitir señales cuando los circuitos de viaje.

Cuadro comparativo

Aproximación a la autolimpieza Ambient Sensing Desacoplado Lógica Personalizada Integración |----------|:-------------:|:---------------:|:---------:|:------------:|:-----------:|

OpenTelemetry # # # # # herramientas externas # # herramientas externas # # herramientas externas # # herramientas externas

Extensiones reactivas Complejo

MediatR # # # # # Manual

Interruptor de circuitos Polly # # # # Interruptor de circuitos Polly # # # Per-call

| Señales efímeras


Señales de elevación y retracción

Ejecución de las operaciones 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; }
}

Dentro de su cuerpo de trabajo:

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

Las señales son sólo cadenas. Utilice nombres simples ("rate-limit") o nombres estructurados ("rate-limit:5000ms"). Los filtros de patrones usan semántica glob (*, ?) y listas de comas de soporte ("error.*,timeout") El emparejamiento es determinista y luz de asignación vía StringPatternMatcher.

Ejemplo de señal de grano fino: llamadas HTTP

Para una observabilidad muy detallada, puede emitir señales en cada etapa de una operación. La biblioteca incluye una muestra SeñalizaciónHttpClient que demuestra este patrón:

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

Esto emite señales en cada etapa:

Señal # # Cuando

|--------|------| | stage.starting Antes de que comience la petición | progress:0 Marca de progreso inicial | stage.request Solicitud HTTP enviada | stage.headers Encabezados de la respuesta recibidos | stage.reading # Comenzando a leer el cuerpo # | progress:XX Porcentaje de progreso (0-100) durante la descarga | stage.completed # Descarga terminada #

A continuación, puede consultar estos con la coincidencia de patrones:

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

Señales retráctiles

Las operaciones también pueden eliminar sus propias señales. Esto es útil para estados temporales:

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

Acontecimientos de retracción

Al igual que la emisión de señal, las retracciones pueden desencadenar callbacks:

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

Los SignalRetractedEvent incluye:

  • Signal - El nombre de la señal retraída
  • OperationId - La operación que la retractó.
  • Key - La llave de la operación (si la hay)
  • Timestamp - Cuando se produjo la retracción
  • WasPatternMatch - Verdadero si se retrae a través de RetractMatching
  • Pattern - El patrón usado (si el patrón coincide)

Ejemplo del mundo real: tasa límite de recuperación

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

Señales de detección

Todos los coordinadores proporcionan una consulta de señal optimizada:

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

Procesamiento de señales reactivas

Desde EfímeroOpciones.cs:

Los coordinadores pueden reaccionar automáticamente a las señales:

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

Cuando una señal en CancelOnSignals se detecta, los nuevos elementos se omiten (contado como fallado). Cuando una señal en DeferOnSignals se detecta, los nuevos elementos esperan hasta que la señal se despeje.


Ejemplo del mundo real: Limitación de la tasa de adaptación

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

Cada instancia de este servicio retrocede automáticamente cuando se golpean los límites de velocidad. No hay estado compartido. No hay mensaje que pase. Sólo se lee la ventana efímera.


Ejemplo del mundo real: Concienciación entre coordinadores

Múltiples coordinadores pueden sentir entre sí a través de un compartido 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);
    }
}

Ejemplo del mundo real: Monitoreo de la Salud

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

No se necesita ninguna biblioteca de métricas. Sólo consulta la ventana efímera.


Ejemplo del mundo real: Rompiendo el circuito basado en señales

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

El disyuntor no tiene su propio estado - sólo lee la ventana efímera.


Limitaciones de la señal: Prevención de los bucles infinitos

Desde Señales.cs:

Cuando las señales pueden causar otras señales, se arriesgan bucles infinitos. SignalConstraints evita esto:

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

Propagación de señales

Causalidad de la vía con 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 cadena de propagación sigue el camino: order-placed → inventory-reserved → ...

Si inventory-reserved tratado de emitir order-placed, sería bloqueado (ciclo detectado).


El fregadero de la señal: el espacio de la señal global

Desde Señales.cs:

Para las señales que deben ser visibles entre los coordinadores:

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

Coincidencia de patrones con StringPatternMatcher

Desde StringPatternMatcher.cs:

Coincidencia al estilo Glob para el filtrado de señales:

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

Convenios sobre la designación de nombres de señales

Mantenga las señales simples y consistentes:

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

Manejo de la señal de ASync

Controladores de señales sincrónicos (OnSignal) ejecutar en el hilo de rosca de la operación, mantenerlos rápidos. Para hacer el trabajo encuadernado con E/S, las señales del ventilador en una ruta asíncrona con SignalDispatcher (combinación de patrones, orden determinista) o AsyncSignalProcessor.

SeñalDispatch fan-out

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

Patrones de soporte *, ?, y listas de comas ("error.*,timeout") Todos los manejadores correspondientes se ejecutan en orden de registro en un coordinador de fondo keyed; emitir sigue siendo sincrónico.

AsyncSignalProcessor

Para el procesamiento independiente async:

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

Ejemplo de mando de señal de telemetría

Un ejemplo completo que combina el procesamiento de señales async con la integración de telemetría:

Fuente: TelemetríaSeñalHandler.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();
}

Encadénelo a su coordinador:

await using var telemetryHandler = new TelemetrySignalHandler(telemetryClient);

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

El encargado:

  • Nunca bloquea el hilo de operación - OnSignal devuelve inmediatamente
  • Señales de categoría por prefijo para diferentes tipos de telemetría
  • Expone métricas (encolado, procesado, eliminado) para el seguimiento de la salud
  • Almacena la memoria - deja caer las señales más antiguas si la cola se llena

¿Por qué funciona esto?

Las señales son poderosas porque son efímero:

Propiedad Beneficio |----------|---------| | Límites No puede crecer sin límite - viejas señales de la edad fuera | Autolimpieza No se necesita código de limpieza | Desconectado Emisores no saben acerca de los oyentes | Observable Cualquier código puede sentir el estado actual | Privado No hay datos de usuario - sólo nombres de señal | Rápido. O(1) detección con cortocircuito

La ventana efímera ya está ahí para la depuración. Las señales sólo le dan significado semántico.


Conclusión

Señales convierten átomos de ejecución aislados en un red de detección. Cada coordinador puede:

  • Emit señales sobre lo que experimentó
  • Sentido señales de su propia historia o de un fregadero compartido
  • Reaccionar automáticamente a través de CancelOnSignals y DeferOnSignals

Ningún corredor de mensajes, ningún estado compartido, ningún protocolo de coordinación, sólo operaciones con metadatos que naturalmente decaen.

Los átomos no hablan entre sí directamente - simplemente dejan rastros en la ventana efímera que otros pueden observar. Es la tigmergia para los sistemas asíncronos.

Fuego... señal... sentido... olvido.


Vínculos

Finding related posts...
logo

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