fix hls direct (#2487)

This commit is contained in:
Jason Dove
2025-10-04 15:36:36 -05:00
committed by GitHub
parent 0363609923
commit c39858b2d8
5 changed files with 137 additions and 91 deletions
+6 -86
View File
@@ -1,6 +1,4 @@
using System.Diagnostics;
using System.IO.Pipelines;
using System.Text;
using CliWrap;
using ErsatzTV.Application.Emby;
using ErsatzTV.Application.Jellyfin;
@@ -23,10 +21,9 @@ namespace ErsatzTV.Controllers;
[ApiController]
[ApiExplorerSettings(IgnoreApi = true)]
public class InternalController : ControllerBase
public class InternalController : StreamingControllerBase
{
private readonly IFFmpegSegmenterService _ffmpegSegmenterService;
private readonly IGraphicsEngine _graphicsEngine;
private readonly ILogger<InternalController> _logger;
private readonly IMediator _mediator;
@@ -35,9 +32,9 @@ public class InternalController : ControllerBase
IGraphicsEngine graphicsEngine,
IMediator mediator,
ILogger<InternalController> logger)
: base(graphicsEngine, logger)
{
_ffmpegSegmenterService = ffmpegSegmenterService;
_graphicsEngine = graphicsEngine;
_mediator = mediator;
_logger = logger;
}
@@ -59,7 +56,7 @@ public class InternalController : ControllerBase
case "segmenter-v2":
return await GetSegmenterV2Stream(channelNumber);
default:
return await GetTsLegacyStream(channelNumber, mode);
return await GetTsLegacyStream(channelNumber);
}
}
@@ -271,14 +268,14 @@ public class InternalController : ControllerBase
worker is HlsSessionWorkerV2 v2)
{
Either<BaseError, PlayoutItemProcessModel> result = await v2.GetNextPlayoutItemProcess();
return GetProcessResponse(result, channelNumber, "segmenter-v2");
return GetProcessResponse(result, channelNumber, StreamingMode.HttpLiveStreamingSegmenterV2);
}
_logger.LogWarning("Unable to locate session worker for channel {Channel}", channelNumber);
return new NotFoundResult();
}
private async Task<IActionResult> GetTsLegacyStream(string channelNumber, string mode)
private async Task<IActionResult> GetTsLegacyStream(string channelNumber)
{
var request = new GetPlayoutItemProcessByChannelNumber(
channelNumber,
@@ -292,83 +289,6 @@ public class InternalController : ControllerBase
Either<BaseError, PlayoutItemProcessModel> result = await _mediator.Send(request);
return GetProcessResponse(result, channelNumber, mode);
}
private IActionResult GetProcessResponse(
Either<BaseError, PlayoutItemProcessModel> result,
string channelNumber,
string mode)
{
foreach (BaseError error in result.LeftToSeq())
{
_logger.LogError(
"Failed to create stream for channel {ChannelNumber}: {Error}",
channelNumber,
error.Value);
return BadRequest(error.Value);
}
foreach (PlayoutItemProcessModel processModel in result.RightToSeq())
{
// for process counter
var ffmpegProcess = new FFmpegProcess();
Command process = processModel.Process;
_logger.LogDebug("ffmpeg arguments {FFmpegArguments}", process.Arguments);
var cts = new CancellationTokenSource();
HttpContext.Response.OnCompleted(async () =>
{
ffmpegProcess.Dispose();
await cts.CancelAsync();
cts.Dispose();
});
using var linkedCts = CancellationTokenSource.CreateLinkedTokenSource(
cts.Token,
HttpContext.RequestAborted);
var pipe = new Pipe();
var stdErrBuffer = new StringBuilder();
Command processWithPipe = process;
foreach (GraphicsEngineContext graphicsEngineContext in processModel.GraphicsEngineContext)
{
var gePipe = new Pipe();
processWithPipe = process.WithStandardInputPipe(PipeSource.FromStream(gePipe.Reader.AsStream()));
// fire and forget graphics engine task
_ = _graphicsEngine.Run(
graphicsEngineContext,
gePipe.Writer,
linkedCts.Token);
}
CommandTask<CommandResult> task = processWithPipe
.WithStandardOutputPipe(PipeTarget.ToStream(pipe.Writer.AsStream()))
.WithStandardErrorPipe(PipeTarget.ToStringBuilder(stdErrBuffer))
.WithValidation(CommandResultValidation.None)
.ExecuteAsync(linkedCts.Token);
// ensure pipe writer is completed when ffmpeg exits
_ = task.Task.ContinueWith(
(_, state) => ((PipeWriter)state!).Complete(),
pipe.Writer,
TaskScheduler.Default);
string contentType = mode switch
{
"segmenter-v2" => "video/x-matroska",
_ => "video/mp2t"
};
return new FileStreamResult(pipe.Reader.AsStream(), contentType);
}
// this will never happen
return new NotFoundResult();
return GetProcessResponse(result, channelNumber, StreamingMode.TransportStream);
}
}
+25 -1
View File
@@ -10,6 +10,7 @@ using ErsatzTV.Core.Domain;
using ErsatzTV.Core.Errors;
using ErsatzTV.Core.FFmpeg;
using ErsatzTV.Core.Interfaces.FFmpeg;
using ErsatzTV.Core.Interfaces.Streaming;
using ErsatzTV.Core.Iptv;
using ErsatzTV.Extensions;
using ErsatzTV.Filters;
@@ -21,7 +22,7 @@ namespace ErsatzTV.Controllers;
[ApiController]
[ApiExplorerSettings(IgnoreApi = true)]
[ServiceFilter(typeof(ConditionalIptvAuthorizeFilter))]
public class IptvController : ControllerBase
public class IptvController : StreamingControllerBase
{
private readonly IFFmpegSegmenterService _ffmpegSegmenterService;
private readonly ILogger<IptvController> _logger;
@@ -29,8 +30,10 @@ public class IptvController : ControllerBase
public IptvController(
IMediator mediator,
IGraphicsEngine graphicsEngine,
ILogger<IptvController> logger,
IFFmpegSegmenterService ffmpegSegmenterService)
: base(graphicsEngine, logger)
{
_mediator = mediator;
_logger = logger;
@@ -284,6 +287,10 @@ public class IptvController : ControllerBase
Right: r => new PhysicalFileResult(r.FileName, r.MimeType));
}
[HttpGet("iptv/hls-direct/{channelNumber}")]
public async Task<IActionResult> GetStream(string channelNumber) =>
await GetHlsDirectStream(channelNumber);
private async Task<string> GetMultiVariantPlaylist(string channelNumber, string mode)
{
string file = mode switch
@@ -359,6 +366,23 @@ public class IptvController : ControllerBase
{variantPlaylist}";
}
private async Task<IActionResult> GetHlsDirectStream(string channelNumber)
{
var request = new GetPlayoutItemProcessByChannelNumber(
channelNumber,
StreamingMode.HttpLiveStreamingDirect,
DateTimeOffset.Now,
false,
true,
DateTimeOffset.Now,
0,
Option<int>.None);
Either<BaseError, PlayoutItemProcessModel> result = await _mediator.Send(request);
return GetProcessResponse(result, channelNumber, StreamingMode.HttpLiveStreamingDirect);
}
private string AccessTokenQuery() => string.IsNullOrWhiteSpace(Request.Query["access_token"])
? string.Empty
: $"?access_token={Request.Query["access_token"]}";
@@ -0,0 +1,92 @@
using System.IO.Pipelines;
using System.Text;
using CliWrap;
using ErsatzTV.Application.Streaming;
using ErsatzTV.Core;
using ErsatzTV.Core.Domain;
using ErsatzTV.Core.FFmpeg;
using ErsatzTV.Core.Interfaces.Streaming;
using Microsoft.AspNetCore.Mvc;
namespace ErsatzTV.Controllers;
public abstract class StreamingControllerBase(IGraphicsEngine graphicsEngine, ILogger logger)
: ControllerBase
{
protected IActionResult GetProcessResponse(
Either<BaseError, PlayoutItemProcessModel> result,
string channelNumber,
StreamingMode mode)
{
foreach (BaseError error in result.LeftToSeq())
{
logger.LogError(
"Failed to create stream for channel {ChannelNumber}: {Error}",
channelNumber,
error.Value);
return BadRequest(error.Value);
}
foreach (PlayoutItemProcessModel processModel in result.RightToSeq())
{
// for process counter
var ffmpegProcess = new FFmpegProcess();
Command process = processModel.Process;
logger.LogDebug("ffmpeg arguments {FFmpegArguments}", process.Arguments);
var cts = new CancellationTokenSource();
HttpContext.Response.OnCompleted(async () =>
{
ffmpegProcess.Dispose();
await cts.CancelAsync();
cts.Dispose();
});
using var linkedCts = CancellationTokenSource.CreateLinkedTokenSource(
cts.Token,
HttpContext.RequestAborted);
var pipe = new Pipe();
var stdErrBuffer = new StringBuilder();
Command processWithPipe = process;
foreach (GraphicsEngineContext graphicsEngineContext in processModel.GraphicsEngineContext)
{
var gePipe = new Pipe();
processWithPipe = process.WithStandardInputPipe(PipeSource.FromStream(gePipe.Reader.AsStream()));
// fire and forget graphics engine task
_ = graphicsEngine.Run(
graphicsEngineContext,
gePipe.Writer,
linkedCts.Token);
}
CommandTask<CommandResult> task = processWithPipe
.WithStandardOutputPipe(PipeTarget.ToStream(pipe.Writer.AsStream()))
.WithStandardErrorPipe(PipeTarget.ToStringBuilder(stdErrBuffer))
.WithValidation(CommandResultValidation.None)
.ExecuteAsync(linkedCts.Token);
// ensure pipe writer is completed when ffmpeg exits
_ = task.Task.ContinueWith(
(_, state) => ((PipeWriter)state!).Complete(),
pipe.Writer,
TaskScheduler.Default);
string contentType = mode switch
{
StreamingMode.HttpLiveStreamingSegmenterV2 => "video/x-matroska",
_ => "video/mp2t"
};
return new FileStreamResult(pipe.Reader.AsStream(), contentType);
}
// this will never happen
return new NotFoundResult();
}
}