Skip to content

Streams y publicación

StreamSource<TItem> publica elementos en memoria. StreamResult<TItem> representa una suscripción y se consume con await foreach.

Productor y consumidor

csharp
using Greencore.Platform.Results;

await using var fuente = new StreamSource<int>();
await using var flujo = StreamResult<int>.Ok(fuente);

var consumidor = Consumir(flujo);
for (var numero = 1; numero <= 3; numero++)
{
    var publicacion = await fuente.PublishAsync(numero);
    Console.WriteLine(publicacion.Status);
}
fuente.Complete();
await consumidor;

static async Task Consumir(StreamResult<int> flujo)
{
    await foreach (var numero in flujo)
        Console.WriteLine($"Recibido: {numero}");
}

Productor y consumidor progresan a la vez. Con un búfer limitado y política Wait, publicar todos los elementos antes de empezar a consumir puede dejar al productor esperando espacio.

Opciones

OpciónPredeterminadoSignificado
ConsumptionModeCompetingConsumersLos enumeradores compiten por los elementos
Capacity256Capacidad; null significa ilimitada
FullModeWaitEspera por espacio
LateEnumeratorModeStartFromCurrentComienza desde el estado actual
ItemDroppedSin callbackNotificación de descartes

SingleConsumer admite un único enumerador activo. Broadcast entrega una copia a cada enumerador activo; no constituye un historial persistente para suscriptores futuros.

Las políticas de saturación disponibles son Wait, Reject, DropWrite, DropOldest y DropNewest. Revisa los descartes mediante DroppedItemCount y ItemDropped; elige una política según si puedes tolerar pérdidas.

Qué confirma una publicación

PublishAsync devuelve StreamPublishResult, con estado Delivered, PartiallyDelivered, Rejected o Canceled y contadores de entregas previstas, aceptadas, rechazadas y canceladas.

La aceptación se refiere al búfer del destino, no a la finalización del procesamiento de negocio. La cancelación de una publicación no retira entregas ya aceptadas ni cancela la suscripción completa.

TryEnqueue es una alternativa inmediata. Un true confirma aceptación en la canalización, no entrega final; un false puede coexistir con aceptaciones parciales entre suscripciones.

Terminación y liberación

Complete() termina normalmente y Fail(exception) termina con error. Una fuente terminada rechaza nuevas publicaciones. La propiedad Termination del flujo informa Completed, Failed o Disposed.

El estado inicial IsOk indica que se creó un resultado correcto; no garantiza que toda la lectura futura termine sin error. Usa await using para liberar fuentes y flujos y atiende las excepciones de lectura cuando corresponda.

Estos streams no ofrecen persistencia, comunicación de red ni garantía de procesamiento exactamente una vez.

Uso personal limitado. Uso comercial sujeto a autorización escrita.