Fix recordings dying mid-stream; record to growing .ts with live buffer playback
🏗️ Build Plugin / build (push) Successful in 12m44s
Nightly Build / nightly-build (push) Successful in 2m49s
🧪 Test Plugin / test (push) Successful in 11m18s
🚀 Release Plugin / build-and-release (push) Successful in 12m24s

A recording of a live event crashed partway through: when its short
Akamai token expired, the proxy's stale-alias check swapped the
recording's mapping to the most recently registered deferred stream,
of any content. Segment requests were then built against the wrong CDN
path, returned 403, and ffmpeg exited, splitting the recording into
several files.

Proxy:
- Never swap URN-backed livestreams; they refresh themselves from the
  API. Only consider swap candidates with the same URN.

Recording lifecycle:
- Complete a recording once its stream has been gone from the API for
  5 minutes (broadcasts can end before ValidTo).
- Fail a recording still waiting for its stream after ValidTo (or 12h
  after ValidFrom) instead of retrying forever.
- Wait for a DVR window whose start= is still in the future instead of
  starting ffmpeg against a 400/404.
- Keep the original start time across restarts; don't mark a recording
  failed on shutdown cancellation.

Growing .ts recordings:
- ffmpeg writes MPEG-TS (-c copy -copyts) to stdout and the plugin
  appends it to one file per recording, so restarts continue the same
  file and timeline.
- The Recordings folder lists any recording with a file on disk, not
  only completed ones.
- In-progress recordings are exposed as infinite streams reading through
  a new endpoint that follows the growing file (same pattern as
  Jellyfin's DVR), so they can be watched from the start like a
  livestream buffer.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
2026-09-26 10:14:24 -04:00
co-authored by Claude Opus 5.5
parent c6a1fd9f63
commit 71a65aab6e
6 changed files with 353 additions and 71 deletions
@@ -81,6 +81,13 @@ public class RecordingEntry
[JsonPropertyName("recordingEndedAt")] [JsonPropertyName("recordingEndedAt")]
public DateTime? RecordingEndedAt { get; set; } public DateTime? RecordingEndedAt { get; set; }
/// <summary>
/// Gets or sets when the stream disappeared from the API after recording had started.
/// Used to detect that a broadcast ended before its scheduled ValidTo.
/// </summary>
[JsonPropertyName("streamLostAt")]
public DateTime? StreamLostAt { get; set; }
/// <summary> /// <summary>
/// Gets or sets the file size in bytes. /// Gets or sets the file size in bytes.
/// </summary> /// </summary>
@@ -98,4 +105,10 @@ public class RecordingEntry
/// </summary> /// </summary>
[JsonPropertyName("createdAt")] [JsonPropertyName("createdAt")]
public DateTime CreatedAt { get; set; } = DateTime.UtcNow; public DateTime CreatedAt { get; set; } = DateTime.UtcNow;
/// <summary>
/// Gets a value indicating whether the output file may still grow (recording, or waiting to resume).
/// </summary>
[JsonIgnore]
public bool IsInProgress => State is RecordingState.Recording or RecordingState.WaitingForStream;
} }
@@ -28,6 +28,8 @@ namespace Jellyfin.Plugin.SRFPlay.Channels;
/// </summary> /// </summary>
public abstract class SrgChannelBase : IChannel, IHasCacheKey public abstract class SrgChannelBase : IChannel, IHasCacheKey
{ {
private const string InProgressPrefix = "● REC · ";
private readonly ILogger<SrgChannelBase> _logger; private readonly ILogger<SrgChannelBase> _logger;
private readonly IContentRefreshService _contentRefreshService; private readonly IContentRefreshService _contentRefreshService;
private readonly IStreamUrlResolver _streamResolver; private readonly IStreamUrlResolver _streamResolver;
@@ -320,15 +322,10 @@ public abstract class SrgChannelBase : IChannel, IHasCacheKey
{ {
var items = new List<ChannelItemInfo>(); var items = new List<ChannelItemInfo>();
var unit = Unit.ToLowerString(); var unit = Unit.ToLowerString();
var recordings = _recordingService.GetRecordings(RecordingState.Completed); var serverBaseUrl = _mediaSourceFactory.GetServerBaseUrl();
foreach (var recording in recordings) foreach (var recording in _recordingService.GetPlayableRecordings())
{ {
if (string.IsNullOrEmpty(recording.OutputPath) || !System.IO.File.Exists(recording.OutputPath))
{
continue;
}
// Only show recordings belonging to this unit. Recordings created before the // Only show recordings belonging to this unit. Recordings created before the
// BusinessUnit tag existed have an empty value and are shown everywhere so they // BusinessUnit tag existed have an empty value and are shown everywhere so they
// are not lost. // are not lost.
@@ -338,43 +335,72 @@ public abstract class SrgChannelBase : IChannel, IHasCacheKey
continue; continue;
} }
var fileInfo = new System.IO.FileInfo(recording.OutputPath); var outputPath = recording.OutputPath!;
var itemId = $"recording_{recording.Id}"; var itemId = $"recording_{recording.Id}";
var container = System.IO.Path.GetExtension(outputPath).TrimStart('.').ToLowerInvariant();
var name = recording.IsInProgress ? $"{InProgressPrefix}{recording.Title}" : recording.Title;
var mediaSource = new MediaSourceInfo MediaSourceInfo mediaSource;
if (recording.IsInProgress)
{
// Same approach as Jellyfin's own DVR: the file is still growing, so treat it as a live
// stream and have ffmpeg read it through an endpoint that follows the file.
mediaSource = new MediaSourceInfo
{ {
Id = itemId, Id = itemId,
Name = recording.Title, Name = name,
Path = recording.OutputPath, Path = outputPath,
Protocol = MediaProtocol.File, Protocol = MediaProtocol.File,
Container = "mkv", EncoderPath = $"{serverBaseUrl}/Plugins/SRFPlay/Recording/{recording.Id}/stream",
EncoderProtocol = MediaProtocol.Http,
Container = container,
SupportsDirectPlay = false,
SupportsDirectStream = true,
SupportsTranscoding = true,
IsInfiniteStream = true,
IgnoreDts = true,
IgnoreIndex = true,
IsRemote = false,
Type = MediaSourceType.Default
};
}
else
{
mediaSource = new MediaSourceInfo
{
Id = itemId,
Name = name,
Path = outputPath,
Protocol = MediaProtocol.File,
Container = container,
SupportsDirectPlay = true, SupportsDirectPlay = true,
SupportsDirectStream = true, SupportsDirectStream = true,
SupportsTranscoding = true, SupportsTranscoding = true,
IsRemote = false, IsRemote = false,
Size = fileInfo.Length, Size = new System.IO.FileInfo(outputPath).Length,
Type = MediaSourceType.Default Type = MediaSourceType.Default
}; };
}
var item = new ChannelItemInfo var item = new ChannelItemInfo
{ {
Id = itemId, Id = itemId,
Name = recording.Title, Name = name,
Overview = recording.Description, Overview = recording.Description,
Type = ChannelItemType.Media, Type = ChannelItemType.Media,
ContentType = ChannelMediaContentType.Movie, ContentType = ChannelMediaContentType.Movie,
MediaType = ChannelMediaType.Video, MediaType = ChannelMediaType.Video,
DateCreated = recording.RecordingStartedAt, DateCreated = recording.RecordingStartedAt,
ImageUrl = !string.IsNullOrEmpty(recording.ImageUrl) ImageUrl = !string.IsNullOrEmpty(recording.ImageUrl)
? CreateProxiedImageUrl(recording.ImageUrl, _mediaSourceFactory.GetServerBaseUrl()) ? CreateProxiedImageUrl(recording.ImageUrl, serverBaseUrl)
: CreatePlaceholderImageUrl(recording.Title, _mediaSourceFactory.GetServerBaseUrl()), : CreatePlaceholderImageUrl(recording.Title, serverBaseUrl),
MediaSources = new List<MediaSourceInfo> { mediaSource } MediaSources = new List<MediaSourceInfo> { mediaSource }
}; };
items.Add(item); items.Add(item);
} }
_logger.LogInformation("Returning {Count} completed recordings as channel items", items.Count); _logger.LogInformation("Returning {Count} recordings as channel items", items.Count);
return items; return items;
} }
@@ -492,8 +518,10 @@ public abstract class SrgChannelBase : IChannel, IHasCacheKey
var timeBucket = new DateTime(now.Year, now.Month, now.Day, now.Hour, (now.Minute / 15) * 15, 0); var timeBucket = new DateTime(now.Year, now.Month, now.Day, now.Hour, (now.Minute / 15) * 15, 0);
var timeKey = timeBucket.ToString("yyyy-MM-dd-HH-mm", CultureInfo.InvariantCulture); var timeKey = timeBucket.ToString("yyyy-MM-dd-HH-mm", CultureInfo.InvariantCulture);
var recordingCount = _recordingService.GetRecordings(RecordingState.Completed).Count; // Include in-progress count so items switch from live to finished sources when a recording ends
return $"{Unit}_{config?.EnableLatestContent}_{config?.EnableTrendingContent}_{config?.EnableCategoryFolders}_{enabledTopics}_{timeKey}_rec{recordingCount}"; var recordings = _recordingService.GetPlayableRecordings();
var recordingKey = $"{recordings.Count}_{recordings.Count(r => r.IsInProgress)}";
return $"{Unit}_{config?.EnableLatestContent}_{config?.EnableTrendingContent}_{config?.EnableCategoryFolders}_{enabledTopics}_{timeKey}_rec{recordingKey}";
} }
private async Task<List<ChannelItemInfo>> ConvertUrnsToChannelItems(List<string> urns, CancellationToken cancellationToken) private async Task<List<ChannelItemInfo>> ConvertUrnsToChannelItems(List<string> urns, CancellationToken cancellationToken)
@@ -1,4 +1,6 @@
using System;
using System.IO; using System.IO;
using System.Linq;
using System.Reflection; using System.Reflection;
using System.Threading; using System.Threading;
using System.Threading.Tasks; using System.Threading.Tasks;
@@ -19,6 +21,9 @@ namespace Jellyfin.Plugin.SRFPlay.Controllers;
[Authorize] [Authorize]
public class RecordingController : ControllerBase public class RecordingController : ControllerBase
{ {
/// <summary>How long a followed stream waits at the end of an in-progress file before giving up.</summary>
private static readonly TimeSpan FollowIdleTimeout = TimeSpan.FromMinutes(2);
private readonly ILogger<RecordingController> _logger; private readonly ILogger<RecordingController> _logger;
private readonly IRecordingService _recordingService; private readonly IRecordingService _recordingService;
@@ -56,6 +61,67 @@ public class RecordingController : ControllerBase
return File(resourceStream, "text/html"); return File(resourceStream, "text/html");
} }
/// <summary>
/// Streams a recording's .ts file. While the recording is still in progress this keeps
/// following the file as it grows, so it can be watched like a livestream buffer.
/// </summary>
/// <param name="id">The recording ID.</param>
/// <param name="cancellationToken">The cancellation token.</param>
/// <returns>The MPEG-TS stream.</returns>
[HttpGet("{id}/stream")]
[AllowAnonymous] // Fetched server-side by Jellyfin/ffmpeg without a token, like the stream proxy
[ProducesResponseType(StatusCodes.Status200OK)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
public async Task<IActionResult> GetRecordingStream(string id, CancellationToken cancellationToken)
{
// The path always comes from our own recording list (which only contains existing files),
// never from the request
var entry = _recordingService.GetPlayableRecordings().FirstOrDefault(r => r.Id == id);
if (entry?.OutputPath is not { } path)
{
return NotFound();
}
Response.ContentType = "video/mp2t";
var input = new FileStream(path, FileMode.Open, FileAccess.Read, FileShare.ReadWrite | FileShare.Delete, 81920, FileOptions.Asynchronous | FileOptions.SequentialScan);
await using (input.ConfigureAwait(false))
{
var buffer = new byte[81920];
var lastDataAt = DateTime.UtcNow;
try
{
while (true)
{
var read = await input.ReadAsync(buffer, cancellationToken).ConfigureAwait(false);
if (read > 0)
{
await Response.Body.WriteAsync(buffer.AsMemory(0, read), cancellationToken).ConfigureAwait(false);
lastDataAt = DateTime.UtcNow;
continue;
}
// At the end of the file: keep following while it's still being written. The idle
// limit covers an ffmpeg restart but ends the stream if recording has stalled.
var stillRecording = _recordingService.GetRecording(id)?.IsInProgress == true;
if (!stillRecording || DateTime.UtcNow - lastDataAt > FollowIdleTimeout)
{
break;
}
await Task.Delay(500, cancellationToken).ConfigureAwait(false);
}
}
catch (OperationCanceledException)
{
// Client went away
}
}
return new EmptyResult();
}
/// <summary> /// <summary>
/// Gets upcoming sport livestreams available for recording. /// Gets upcoming sport livestreams available for recording.
/// </summary> /// </summary>
@@ -47,6 +47,19 @@ public interface IRecordingService
/// <returns>List of matching recording entries.</returns> /// <returns>List of matching recording entries.</returns>
IReadOnlyList<RecordingEntry> GetRecordings(RecordingState? stateFilter = null); IReadOnlyList<RecordingEntry> GetRecordings(RecordingState? stateFilter = null);
/// <summary>
/// Gets a single recording by ID.
/// </summary>
/// <param name="recordingId">The recording ID.</param>
/// <returns>The recording entry, or null if not found.</returns>
RecordingEntry? GetRecording(string recordingId);
/// <summary>
/// Gets recordings that have an output file on disk, finished or still in progress.
/// </summary>
/// <returns>Playable recordings, newest first.</returns>
IReadOnlyList<RecordingEntry> GetPlayableRecordings();
/// <summary> /// <summary>
/// Deletes a completed recording (entry and optionally the file). /// Deletes a completed recording (entry and optionally the file).
/// </summary> /// </summary>
@@ -25,6 +25,15 @@ namespace Jellyfin.Plugin.SRFPlay.Services;
/// </summary> /// </summary>
public class RecordingService : IRecordingService, IDisposable public class RecordingService : IRecordingService, IDisposable
{ {
/// <summary>How long the stream may be missing from the API mid-recording before we treat the broadcast as over.</summary>
private static readonly TimeSpan StreamLostGracePeriod = TimeSpan.FromMinutes(5);
/// <summary>When a recording has no ValidTo, how long after ValidFrom to keep waiting for the stream.</summary>
private static readonly TimeSpan NoValidToGiveUpAfter = TimeSpan.FromHours(12);
/// <summary>Longest the scheduler will block waiting for an imminent DVR window start.</summary>
private static readonly TimeSpan MaxInlineStartWait = TimeSpan.FromSeconds(60);
private readonly ILogger<RecordingService> _logger; private readonly ILogger<RecordingService> _logger;
private readonly ISRFApiClientFactory _apiClientFactory; private readonly ISRFApiClientFactory _apiClientFactory;
private readonly IStreamProxyService _proxyService; private readonly IStreamProxyService _proxyService;
@@ -32,7 +41,7 @@ public class RecordingService : IRecordingService, IDisposable
private readonly IMediaCompositionFetcher _mediaCompositionFetcher; private readonly IMediaCompositionFetcher _mediaCompositionFetcher;
private readonly IServerApplicationHost _appHost; private readonly IServerApplicationHost _appHost;
private readonly IMediaEncoder _mediaEncoder; private readonly IMediaEncoder _mediaEncoder;
private readonly ConcurrentDictionary<string, Process> _activeProcesses = new(); private readonly ConcurrentDictionary<string, ActiveRecording> _activeProcesses = new();
private static readonly JsonSerializerOptions _jsonOptions = new() { WriteIndented = true }; private static readonly JsonSerializerOptions _jsonOptions = new() { WriteIndented = true };
private readonly SemaphoreSlim _persistLock = new(1, 1); private readonly SemaphoreSlim _persistLock = new(1, 1);
private readonly SemaphoreSlim _processLock = new(1, 1); private readonly SemaphoreSlim _processLock = new(1, 1);
@@ -266,15 +275,7 @@ public class RecordingService : IRecordingService, IDisposable
} }
StopFfmpeg(recordingId); StopFfmpeg(recordingId);
CompleteRecording(entry, DateTime.UtcNow);
entry.State = RecordingState.Completed;
entry.RecordingEndedAt = DateTime.UtcNow;
if (entry.OutputPath != null && File.Exists(entry.OutputPath))
{
entry.FileSizeBytes = new FileInfo(entry.OutputPath).Length;
}
_ = SaveRecordingsAsync(); _ = SaveRecordingsAsync();
_logger.LogInformation("Stopped recording '{Title}' ({Id})", entry.Title, recordingId); _logger.LogInformation("Stopped recording '{Title}' ({Id})", entry.Title, recordingId);
@@ -298,6 +299,22 @@ public class RecordingService : IRecordingService, IDisposable
return _recordings.OrderByDescending(r => r.CreatedAt).ToList(); return _recordings.OrderByDescending(r => r.CreatedAt).ToList();
} }
/// <inheritdoc />
public RecordingEntry? GetRecording(string recordingId)
{
return GetRecordings(null).FirstOrDefault(r => r.Id == recordingId);
}
/// <inheritdoc />
public IReadOnlyList<RecordingEntry> GetPlayableRecordings()
{
// Gate on the file rather than the state: an in-progress or failed recording
// still has watchable footage.
return GetRecordings(null)
.Where(r => !string.IsNullOrEmpty(r.OutputPath) && File.Exists(r.OutputPath))
.ToList();
}
/// <inheritdoc /> /// <inheritdoc />
public bool DeleteRecording(string recordingId, bool deleteFile) public bool DeleteRecording(string recordingId, bool deleteFile)
{ {
@@ -369,8 +386,26 @@ public class RecordingService : IRecordingService, IDisposable
{ {
case RecordingState.Scheduled: case RecordingState.Scheduled:
case RecordingState.WaitingForStream: case RecordingState.WaitingForStream:
// Check if it's time to start recording // Give up once the broadcast window is over, otherwise we'd retry forever
if (validFromUtc.HasValue && validFromUtc.Value <= now.AddMinutes(2)) var giveUpAt = validToUtc ?? validFromUtc?.Add(NoValidToGiveUpAfter);
if (giveUpAt.HasValue && giveUpAt.Value <= now)
{
if (entry.RecordingStartedAt.HasValue)
{
_logger.LogInformation("Recording '{Title}' reached end of broadcast window while waiting for stream, completing", entry.Title);
CompleteRecording(entry, now);
}
else
{
_logger.LogWarning("Stream for '{Title}' never became available (window ended {GiveUpAt}), marking failed", entry.Title, giveUpAt.Value);
entry.State = RecordingState.Failed;
entry.ErrorMessage = "Stream never became available";
entry.RecordingEndedAt = now;
}
changed = true;
}
else if (validFromUtc.HasValue && validFromUtc.Value <= now.AddMinutes(2))
{ {
_logger.LogInformation( _logger.LogInformation(
"Time to start recording '{Title}': ValidFrom={ValidFrom} (UTC: {ValidFromUtc}), Now={Now}", "Time to start recording '{Title}': ValidFrom={ValidFrom} (UTC: {ValidFromUtc}), Now={Now}",
@@ -389,13 +424,7 @@ public class RecordingService : IRecordingService, IDisposable
{ {
_logger.LogInformation("Recording '{Title}' reached ValidTo, stopping", entry.Title); _logger.LogInformation("Recording '{Title}' reached ValidTo, stopping", entry.Title);
StopFfmpeg(entry.Id); StopFfmpeg(entry.Id);
entry.State = RecordingState.Completed; CompleteRecording(entry, now);
entry.RecordingEndedAt = now;
if (entry.OutputPath != null && File.Exists(entry.OutputPath))
{
entry.FileSizeBytes = new FileInfo(entry.OutputPath).Length;
}
changed = true; changed = true;
} }
else if (!_activeProcesses.ContainsKey(entry.Id)) else if (!_activeProcesses.ContainsKey(entry.Id))
@@ -426,7 +455,7 @@ public class RecordingService : IRecordingService, IDisposable
if (chapter == null) if (chapter == null)
{ {
_logger.LogDebug("No chapter found for '{Title}', stream may not be live yet", entry.Title); _logger.LogDebug("No chapter found for '{Title}', stream may not be live yet", entry.Title);
entry.State = RecordingState.WaitingForStream; MarkStreamUnavailable(entry);
return true; return true;
} }
@@ -437,10 +466,29 @@ public class RecordingService : IRecordingService, IDisposable
if (string.IsNullOrEmpty(streamUrl)) if (string.IsNullOrEmpty(streamUrl))
{ {
_logger.LogDebug("No stream URL available for '{Title}', waiting", entry.Title); _logger.LogDebug("No stream URL available for '{Title}', waiting", entry.Title);
MarkStreamUnavailable(entry);
return true;
}
// The CDN rejects a DVR window whose start= is still in the future (400/404).
// Wait briefly if it's imminent, otherwise try again on the next scheduler run.
var windowStart = GetDvrWindowStart(streamUrl);
if (windowStart.HasValue && windowStart.Value > DateTime.UtcNow)
{
var wait = windowStart.Value - DateTime.UtcNow + TimeSpan.FromSeconds(5);
if (wait > MaxInlineStartWait)
{
_logger.LogInformation("Stream for '{Title}' starts at {Start}, waiting", entry.Title, windowStart.Value);
entry.State = RecordingState.WaitingForStream; entry.State = RecordingState.WaitingForStream;
return true; return true;
} }
_logger.LogInformation("Stream for '{Title}' starts in {Seconds:F0}s, delaying ffmpeg start", entry.Title, wait.TotalSeconds);
await Task.Delay(wait, cancellationToken).ConfigureAwait(false);
}
entry.StreamLostAt = null;
// Register the stream with the proxy so we can use the proxy URL // Register the stream with the proxy so we can use the proxy URL
var itemId = $"rec_{entry.Id}"; var itemId = $"rec_{entry.Id}";
var isLiveStream = chapter.Type == "SCHEDULED_LIVESTREAM" || UrnHelper.IsLivestreamUrn(entry.Urn); var isLiveStream = chapter.Type == "SCHEDULED_LIVESTREAM" || UrnHelper.IsLivestreamUrn(entry.Urn);
@@ -449,22 +497,26 @@ public class RecordingService : IRecordingService, IDisposable
// Build proxy URL for ffmpeg (use localhost for local access) // Build proxy URL for ffmpeg (use localhost for local access)
var proxyUrl = $"{GetServerBaseUrl()}/Plugins/SRFPlay/Proxy/{itemId}/master.m3u8"; var proxyUrl = $"{GetServerBaseUrl()}/Plugins/SRFPlay/Proxy/{itemId}/master.m3u8";
// Build output file path // Restarts append to the existing .ts so a recording stays one file
var outputPath = entry.OutputPath;
if (string.IsNullOrEmpty(outputPath) || !outputPath.EndsWith(".ts", StringComparison.OrdinalIgnoreCase))
{
var safeTitle = SanitizeFileName(entry.Title); var safeTitle = SanitizeFileName(entry.Title);
var timestamp = DisplayTime.Now().ToString("yyyy-MM-dd_HHmm", CultureInfo.InvariantCulture); var timestamp = DisplayTime.Now().ToString("yyyy-MM-dd_HHmm", CultureInfo.InvariantCulture);
var outputPath = Path.Combine(GetRecordingOutputPath(), $"{safeTitle}_{timestamp}.mkv"); outputPath = Path.Combine(GetRecordingOutputPath(), $"{safeTitle}_{timestamp}.ts");
entry.OutputPath = outputPath; entry.OutputPath = outputPath;
}
// Start ffmpeg // Start ffmpeg
StartFfmpeg(entry.Id, proxyUrl, outputPath); StartFfmpeg(entry.Id, proxyUrl, outputPath);
entry.State = RecordingState.Recording; entry.State = RecordingState.Recording;
entry.RecordingStartedAt = DateTime.UtcNow; entry.RecordingStartedAt ??= DateTime.UtcNow;
_logger.LogInformation("Started recording '{Title}' to {OutputPath}", entry.Title, outputPath); _logger.LogInformation("Started recording '{Title}' to {OutputPath}", entry.Title, outputPath);
return true; return true;
} }
catch (Exception ex) catch (Exception ex) when (ex is not OperationCanceledException)
{ {
_logger.LogError(ex, "Failed to start recording '{Title}'", entry.Title); _logger.LogError(ex, "Failed to start recording '{Title}'", entry.Title);
entry.State = RecordingState.Failed; entry.State = RecordingState.Failed;
@@ -473,6 +525,61 @@ public class RecordingService : IRecordingService, IDisposable
} }
} }
/// <summary>
/// Handles the API no longer offering a stream. Before recording starts this just means "not live yet";
/// once recording has started it means the broadcast ended (possibly before ValidTo).
/// </summary>
private void MarkStreamUnavailable(RecordingEntry entry)
{
if (!entry.RecordingStartedAt.HasValue)
{
entry.State = RecordingState.WaitingForStream;
return;
}
var now = DateTime.UtcNow;
entry.StreamLostAt ??= now;
if (now - entry.StreamLostAt.Value >= StreamLostGracePeriod)
{
_logger.LogInformation(
"Stream for '{Title}' has been gone since {LostAt}, broadcast has ended - completing recording",
entry.Title,
entry.StreamLostAt.Value);
StopFfmpeg(entry.Id);
CompleteRecording(entry, now);
return;
}
entry.State = RecordingState.WaitingForStream;
}
private static void CompleteRecording(RecordingEntry entry, DateTime now)
{
entry.State = RecordingState.Completed;
entry.RecordingEndedAt = now;
if (entry.OutputPath != null && File.Exists(entry.OutputPath))
{
entry.FileSizeBytes = new FileInfo(entry.OutputPath).Length;
}
}
/// <summary>
/// Reads the DVR window start (unix seconds in the start= query parameter) from a stream URL.
/// </summary>
private static DateTime? GetDvrWindowStart(string streamUrl)
{
if (!Uri.TryCreate(streamUrl, UriKind.Absolute, out var uri))
{
return null;
}
var start = System.Web.HttpUtility.ParseQueryString(uri.Query)["start"];
return long.TryParse(start, NumberStyles.Integer, CultureInfo.InvariantCulture, out var seconds)
? DateTimeOffset.FromUnixTimeSeconds(seconds).UtcDateTime
: null;
}
private void StartFfmpeg(string recordingId, string inputUrl, string outputPath) private void StartFfmpeg(string recordingId, string inputUrl, string outputPath)
{ {
var process = new Process var process = new Process
@@ -480,13 +587,16 @@ public class RecordingService : IRecordingService, IDisposable
StartInfo = new ProcessStartInfo StartInfo = new ProcessStartInfo
{ {
FileName = _mediaEncoder.EncoderPath, FileName = _mediaEncoder.EncoderPath,
Arguments = $"-y -i \"{inputUrl}\" -c copy -movflags +faststart \"{outputPath}\"", // MPEG-TS on stdout: it has no header/index to finalise, so the file is playable while it
// grows, and we append it ourselves so restarts continue the same file. -copyts keeps the
// broadcast timestamps so an appended restart continues the timeline instead of resetting to 0.
Arguments = $"-i \"{inputUrl}\" -c copy -copyts -f mpegts pipe:1",
UseShellExecute = false, UseShellExecute = false,
RedirectStandardInput = true, RedirectStandardInput = true,
RedirectStandardOutput = true,
RedirectStandardError = true, RedirectStandardError = true,
CreateNoWindow = true CreateNoWindow = true
}, }
EnableRaisingEvents = true
}; };
process.ErrorDataReceived += (_, args) => process.ErrorDataReceived += (_, args) =>
@@ -497,23 +607,49 @@ public class RecordingService : IRecordingService, IDisposable
} }
}; };
process.Exited += (_, _) =>
{
_logger.LogInformation("ffmpeg process exited for recording {RecordingId} with code {ExitCode}", recordingId, process.ExitCode);
_activeProcesses.TryRemove(recordingId, out _);
};
process.Start(); process.Start();
process.BeginErrorReadLine(); process.BeginErrorReadLine();
_activeProcesses[recordingId] = process; var active = new ActiveRecording(process);
_logger.LogInformation("Started ffmpeg (PID {Pid}) for recording {RecordingId}: {Args}", process.Id, recordingId, process.StartInfo.Arguments); _activeProcesses[recordingId] = active;
active.Pump = PumpToFileAsync(recordingId, active, outputPath);
_logger.LogInformation("Started ffmpeg (PID {Pid}) for recording {RecordingId} -> {OutputPath}: {Args}", process.Id, recordingId, outputPath, process.StartInfo.Arguments);
}
/// <summary>
/// Appends ffmpeg's stdout to the recording file. The recording counts as active until
/// ffmpeg has exited and all of its output is on disk, so a restart never interleaves writes.
/// </summary>
private async Task PumpToFileAsync(string recordingId, ActiveRecording active, string outputPath)
{
try
{
var output = new FileStream(outputPath, FileMode.Append, FileAccess.Write, FileShare.ReadWrite | FileShare.Delete);
await using (output.ConfigureAwait(false))
{
await active.Process.StandardOutput.BaseStream.CopyToAsync(output).ConfigureAwait(false);
}
await active.Process.WaitForExitAsync().ConfigureAwait(false);
_logger.LogInformation("ffmpeg process exited for recording {RecordingId} with code {ExitCode}", recordingId, active.Process.ExitCode);
}
catch (Exception ex)
{
_logger.LogWarning(ex, "Error writing recording output for {RecordingId}", recordingId);
}
finally
{
// Only remove our own entry; a newer process may already be registered
_activeProcesses.TryRemove(new KeyValuePair<string, ActiveRecording>(recordingId, active));
}
} }
private void StopFfmpeg(string recordingId) private void StopFfmpeg(string recordingId)
{ {
if (_activeProcesses.TryRemove(recordingId, out var process)) if (_activeProcesses.TryRemove(recordingId, out var active))
{ {
var process = active.Process;
try try
{ {
if (!process.HasExited) if (!process.HasExited)
@@ -529,6 +665,8 @@ public class RecordingService : IRecordingService, IDisposable
} }
} }
// Let the remaining output reach the file before releasing the process
active.Pump?.Wait(TimeSpan.FromSeconds(10));
process.Dispose(); process.Dispose();
} }
catch (Exception ex) catch (Exception ex)
@@ -578,4 +716,19 @@ public class RecordingService : IRecordingService, IDisposable
_disposed = true; _disposed = true;
} }
/// <summary>
/// A running ffmpeg process and the task copying its output to disk.
/// </summary>
private sealed class ActiveRecording
{
public ActiveRecording(Process process)
{
Process = process;
}
public Process Process { get; }
public Task? Pump { get; set; }
}
} }
@@ -132,10 +132,13 @@ public class StreamProxyService : IStreamProxyService
!string.IsNullOrEmpty(streamInfo.AuthenticatedUrl)); !string.IsNullOrEmpty(streamInfo.AuthenticatedUrl));
// Check for stale alias: only look for fresher stream if current token is EXPIRED or EXPIRING SOON // Check for stale alias: only look for fresher stream if current token is EXPIRED or EXPIRING SOON
// Don't replace a valid token (>5s left) with a new deferred registration // Don't replace a valid token (>5s left) with a new deferred registration.
if (!streamInfo.NeedsAuthentication && tokenTimeLeft < 5) // Livestreams with a URN refresh themselves from the API, so never swap them: an unrelated
// registration (e.g. from a library refresh) would point the stream at the wrong CDN path.
var selfRefreshing = streamInfo.IsLiveStream && !string.IsNullOrEmpty(streamInfo.Urn);
if (!streamInfo.NeedsAuthentication && tokenTimeLeft < 5 && !selfRefreshing)
{ {
var freshStream = FindFreshestStream(); var freshStream = FindFreshestStream(streamInfo.Urn);
if (freshStream != null && freshStream.Value.Value.NeedsAuthentication) if (freshStream != null && freshStream.Value.Value.NeedsAuthentication)
{ {
_logger.LogWarning( _logger.LogWarning(
@@ -595,8 +598,9 @@ public class StreamProxyService : IStreamProxyService
/// <summary> /// <summary>
/// Finds the freshest (most recently registered) stream that needs authentication or has a valid token. /// Finds the freshest (most recently registered) stream that needs authentication or has a valid token.
/// </summary> /// </summary>
/// <param name="urn">When set, only streams for this URN are considered.</param>
/// <returns>The freshest stream entry, or null if none found.</returns> /// <returns>The freshest stream entry, or null if none found.</returns>
private KeyValuePair<string, StreamInfo>? FindFreshestStream() private KeyValuePair<string, StreamInfo>? FindFreshestStream(string? urn)
{ {
var now = DateTime.UtcNow; var now = DateTime.UtcNow;
@@ -604,6 +608,11 @@ public class StreamProxyService : IStreamProxyService
// or have tokens that aren't expired yet // or have tokens that aren't expired yet
var candidates = _streamMappings.Where(kvp => var candidates = _streamMappings.Where(kvp =>
{ {
if (!string.IsNullOrEmpty(urn) && !string.Equals(kvp.Value.Urn, urn, StringComparison.OrdinalIgnoreCase))
{
return false; // Different content
}
if (kvp.Value.NeedsAuthentication) if (kvp.Value.NeedsAuthentication)
{ {
return true; // Fresh deferred registration return true; // Fresh deferred registration