Greencore.Platform.Streaming
0.0.1-alpha.2
dotnet add package Greencore.Platform.Streaming --version 0.0.1-alpha.2
NuGet\Install-Package Greencore.Platform.Streaming -Version 0.0.1-alpha.2
<PackageReference Include="Greencore.Platform.Streaming" Version="0.0.1-alpha.2" />
<PackageVersion Include="Greencore.Platform.Streaming" Version="0.0.1-alpha.2" />
<PackageReference Include="Greencore.Platform.Streaming" />
paket add Greencore.Platform.Streaming --version 0.0.1-alpha.2
#r "nuget: Greencore.Platform.Streaming, 0.0.1-alpha.2"
#:package Greencore.Platform.Streaming@0.0.1-alpha.2
#addin nuget:?package=Greencore.Platform.Streaming&version=0.0.1-alpha.2&prerelease
#tool nuget:?package=Greencore.Platform.Streaming&version=0.0.1-alpha.2&prerelease
Greencore.platform.streaming
Publicación en memoria mediante Channel<T>. El núcleo solo depende de .NET: no referencia Results, Errors, Localization ni Concurrency.
using Greencore.platform.streaming;
await using var hub = StreamHub<int>.CreateBroadcast();
await using var first = hub.Subscribe();
await using var second = hub.Subscribe();
var firstTask = ConsumeAsync(first);
var secondTask = ConsumeAsync(second);
var outcome = await hub.PublishAsync(42);
hub.Complete();
await Task.WhenAll(firstTask, secondTask);
static async Task ConsumeAsync(StreamSubscription<int> subscription)
{
await foreach (var item in subscription.ReadAllAsync())
Console.WriteLine(item);
}
Distribución y capacidad
CreateBroadcast() crea un canal por suscripción. Los destinatarios se capturan para cada publicación; una suscripción nueva no recibe el historial. Sin suscripciones, PublishAsync devuelve NoRecipients.
CreateWorkQueue() crea un canal compartido. Puede almacenar elementos sin consumidores hasta alcanzar capacidad; cada elemento se retira una sola vez. No hay confirmación de procesamiento ni reintento automático si un consumidor falla después de retirarlo.
StreamHubOptions.Capacity debe ser positiva. FullMode utiliza BoundedChannelFullMode: Wait aplica backpressure, mientras DropWrite, DropOldest y DropNewest descartan según la política de .NET. En broadcast un consumidor lento puede retrasar la publicación. No se garantiza equidad entre consumidores competidores.
Las publicaciones se serializan: todos los destinatarios que reciben los mismos elementos observan el mismo orden. No se promete un orden específico entre llamadas simultáneas. El búfer es acotado; el productor también debe limitar la cantidad de llamadas simultáneas a PublishAsync si necesita acotar las publicaciones pendientes.
StreamPublishOutcome.Accepted confirma aceptación en el canal, no procesamiento. Rejected representa canales cerrados; Canceled representa entregas interrumpidas. Dropped cuenta descartes durante esa publicación y DroppedItemCount acumula los descartes por capacidad. En DropWrite el elemento entrante no cuenta como aceptado; en los otros modos se acepta mientras se descarta uno anterior. Los descartes se cuentan por canal, no como elementos únicos globales. El descarte del búfer al disponer no incrementa este contador.
Si la cancelación del consumidor ocurre antes de capturar destinatarios, el resultado tiene Recipients = 0 y CanceledBeforeDelivery = true. Closed indica que el hub no admite la publicación.
Consumo y cancelación
Una suscripción permite un solo enumerador. Para varios consumidores llama varias veces a Subscribe(). Un segundo MoveNextAsync pendiente se rechaza. Current no debe utilizarse concurrentemente con movimiento o disposición.
Los tokens pasados a Subscribe(token), ReadAllAsync(token) y GetAsyncEnumerator(token) cancelan únicamente esa suscripción. Todas las rutas participan en la misma limpieza. CancelAsync() solicita cancelación y espera la limpieza de la suscripción. En broadcast se descarta su cola; en work queue permanecen los elementos que todavía no fueron retirados.
El token de PublishAsync(item, token) cancela solamente esa publicación. Las entregas ya aceptadas se conservan y una publicación broadcast puede ser parcial.
Complete() impide publicar y permite drenar canales. Fail(exception) conserva los elementos aceptados y propaga la excepción al terminar el drenaje. La primera terminación del productor prevalece. No se permiten nuevas suscripciones después del cierre. Completion describe la terminación de cada suscripción después de limpiar sus recursos; una suscripción que nunca se enumera debe cancelarse o disponerse explícitamente.
DisposeAsync() del hub aborta las suscripciones, desbloquea publicaciones pendientes y espera la limpieza interna. Todas las llamadas comparten una única tarea de disposición. Solo se ofrece disposición asíncrona.
La disposición no espera el trabajo arbitrario dentro del cuerpo de await foreach. Para esperar también ese trabajo, conserva sus tareas, pasa un token al procesamiento y espera Task.WhenAll(consumers) después de detener el hub.
Migración
StreamSource<T>pasa aStreamHub<T>.StreamResult<T>.Valuespasa ahub.Subscribe().ReadAllAsync().TryEnqueuese sustituye porPublishAsync, con un resultado explícito de entregas parciales.- El modo de consumidor único se obtiene creando una sola suscripción.
- No se conservan capacidad ilimitada,
Reject, callbacks de descarte ni rechazo de consumidores tardíos como opciones de la nueva API. StreamResult<T>se ofrece en el proyecto opcionalplatform.results.streaming, cuyo valor es el hub. No es un consumidor ni dispone el hub.
Se trata de una migración incompatible de la API anterior. La implementación antigua permanece recuperable en el historial de Git.
| Product | Versions Compatible and additional computed target framework versions. |
|---|---|
| .NET | net10.0 is compatible. net10.0-android was computed. net10.0-browser was computed. net10.0-ios was computed. net10.0-maccatalyst was computed. net10.0-macos was computed. net10.0-tvos was computed. net10.0-windows was computed. |
-
net10.0
- No dependencies.
NuGet packages (1)
Showing the top 1 NuGet packages that depend on Greencore.Platform.Streaming:
| Package | Downloads |
|---|---|
|
Greencore.Platform.Results.Streaming
Resultados explícitos, errores estructurados, validación y concurrencia asíncrona para .NET. |
GitHub repositories
This package is not used by any popular GitHub repositories.
| Version | Downloads | Last Updated |
|---|---|---|
| 0.0.1-alpha.2 | 41 | 9/26/2026 |
Siguiente versión alpha.2 de Greencore.