dragon-iso/src/TeamsISO.Engine/Pipeline/NdiSender.cs

71 lines
2.3 KiB
C#
Raw Normal View History

using System.Threading.Channels;
using Microsoft.Extensions.Logging;
using TeamsISO.Engine.Interop;
namespace TeamsISO.Engine.Pipeline;
/// <summary>
/// Pulls processed frames from a channel and forwards them to <see cref="INdiInterop.SendFrame"/>.
/// </summary>
public sealed class NdiSender : IDisposable
{
private readonly INdiInterop _interop;
private readonly string _outputName;
private readonly ChannelReader<ProcessedFrame> _input;
private readonly ILogger<NdiSender> _logger;
private readonly NdiSenderHandle _handle;
private long _framesSent;
public NdiSender(
INdiInterop interop,
string outputName,
ChannelReader<ProcessedFrame> input,
ILogger<NdiSender> logger)
{
_interop = interop;
_outputName = outputName;
_input = input;
_logger = logger;
_handle = interop.CreateSender(outputName);
}
public long FramesSent => Interlocked.Read(ref _framesSent);
/// <summary>
/// Awaits one frame and forwards it. Returns false if the channel is completed.
/// Test seam.
/// </summary>
public async ValueTask<bool> SendNextAsync(CancellationToken cancellationToken)
{
if (!await _input.WaitToReadAsync(cancellationToken))
return false;
if (!_input.TryRead(out var frame))
return false;
_interop.SendFrame(_handle, frame);
Interlocked.Increment(ref _framesSent);
return true;
}
/// <summary>Long-running send loop. Run on a dedicated thread.</summary>
public Task RunAsync(CancellationToken cancellationToken) =>
Task.Factory.StartNew(async () =>
{
try
{
while (!cancellationToken.IsCancellationRequested)
{
var more = await SendNextAsync(cancellationToken);
if (!more) break;
}
}
catch (OperationCanceledException) { }
catch (Exception ex)
{
_logger.LogError(ex, "NdiSender loop crashed for output {Output}.", _outputName);
throw;
}
}, cancellationToken, TaskCreationOptions.LongRunning, TaskScheduler.Default).Unwrap();
public void Dispose() => _handle.Dispose();
}