using System;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Diagnostics;
using System.Globalization;
using System.IO;
using System.Linq;
using System.Text.Json;
using System.Text.RegularExpressions;
using System.Threading;
using System.Threading.Tasks;
using Jellyfin.Plugin.SRFPlay.Api;
using Jellyfin.Plugin.SRFPlay.Api.Models;
using Jellyfin.Plugin.SRFPlay.Api.Models.PlayV3;
using Jellyfin.Plugin.SRFPlay.Services.Interfaces;
using Jellyfin.Plugin.SRFPlay.Utilities;
using MediaBrowser.Controller;
using MediaBrowser.Controller.MediaEncoding;
using Microsoft.Extensions.Logging;
namespace Jellyfin.Plugin.SRFPlay.Services;
///
/// Service for managing sport livestream recordings using ffmpeg.
///
public class RecordingService : IRecordingService, IDisposable
{
/// How long the stream may be missing from the API mid-recording before we treat the broadcast as over.
private static readonly TimeSpan StreamLostGracePeriod = TimeSpan.FromMinutes(5);
/// When a recording has no ValidTo, how long after ValidFrom to keep waiting for the stream.
private static readonly TimeSpan NoValidToGiveUpAfter = TimeSpan.FromHours(12);
/// Longest the scheduler will block waiting for an imminent DVR window start.
private static readonly TimeSpan MaxInlineStartWait = TimeSpan.FromSeconds(60);
private readonly ILogger _logger;
private readonly ISRFApiClientFactory _apiClientFactory;
private readonly IStreamProxyService _proxyService;
private readonly IStreamUrlResolver _streamUrlResolver;
private readonly IMediaCompositionFetcher _mediaCompositionFetcher;
private readonly IServerApplicationHost _appHost;
private readonly IMediaEncoder _mediaEncoder;
private readonly ConcurrentDictionary _activeProcesses = new();
private static readonly JsonSerializerOptions _jsonOptions = new() { WriteIndented = true };
private readonly SemaphoreSlim _persistLock = new(1, 1);
private readonly SemaphoreSlim _processLock = new(1, 1);
private List _recordings = new();
private bool _loaded;
private bool _disposed;
///
/// Initializes a new instance of the class.
///
/// The logger.
/// The API client factory.
/// The stream proxy service.
/// The stream URL resolver.
/// The media composition fetcher.
/// The application host.
/// The media encoder for ffmpeg path.
public RecordingService(
ILogger logger,
ISRFApiClientFactory apiClientFactory,
IStreamProxyService proxyService,
IStreamUrlResolver streamUrlResolver,
IMediaCompositionFetcher mediaCompositionFetcher,
IServerApplicationHost appHost,
IMediaEncoder mediaEncoder)
{
_logger = logger;
_apiClientFactory = apiClientFactory;
_proxyService = proxyService;
_streamUrlResolver = streamUrlResolver;
_mediaCompositionFetcher = mediaCompositionFetcher;
_appHost = appHost;
_mediaEncoder = mediaEncoder;
}
private string GetDataFilePath()
{
var dataPath = Plugin.Instance?.DataFolderPath ?? Path.Combine(Environment.GetFolderPath(Environment.SpecialFolder.ApplicationData), "jellyfin", "plugins", "SRFPlay");
Directory.CreateDirectory(dataPath);
return Path.Combine(dataPath, "recordings.json");
}
private string GetRecordingOutputPath()
{
var config = Plugin.Instance?.Configuration;
var path = config?.RecordingOutputPath;
if (string.IsNullOrWhiteSpace(path))
{
path = Path.Combine(Environment.GetFolderPath(Environment.SpecialFolder.UserProfile), "SRFRecordings");
}
Directory.CreateDirectory(path);
return path;
}
private string GetServerBaseUrl()
{
var config = Plugin.Instance?.Configuration;
if (config != null && !string.IsNullOrWhiteSpace(config.PublicServerUrl))
{
return config.PublicServerUrl.TrimEnd('/');
}
// For local ffmpeg access, use localhost directly
return "http://localhost:8096";
}
private async Task LoadRecordingsAsync()
{
if (_loaded)
{
return;
}
var filePath = GetDataFilePath();
if (File.Exists(filePath))
{
try
{
var json = await File.ReadAllTextAsync(filePath).ConfigureAwait(false);
_recordings = JsonSerializer.Deserialize>(json) ?? new List();
_logger.LogInformation("Loaded {Count} recording entries from {Path}", _recordings.Count, filePath);
}
catch (Exception ex)
{
_logger.LogError(ex, "Failed to load recordings from {Path}", filePath);
_recordings = new List();
}
}
_loaded = true;
}
private async Task SaveRecordingsAsync()
{
await _persistLock.WaitAsync().ConfigureAwait(false);
try
{
var filePath = GetDataFilePath();
var json = JsonSerializer.Serialize(_recordings, _jsonOptions);
await File.WriteAllTextAsync(filePath, json).ConfigureAwait(false);
}
catch (Exception ex)
{
_logger.LogError(ex, "Failed to save recordings");
}
finally
{
_persistLock.Release();
}
}
///
public async Task> GetUpcomingScheduleAsync(CancellationToken cancellationToken)
{
var units = System.Enum.GetValues();
using var apiClient = _apiClientFactory.CreateClient();
// Aggregate sport livestreams across every business unit so the recordings
// page shows events from all languages.
var all = new List();
foreach (var unit in units)
{
var businessUnit = unit.ToString().ToLowerInvariant();
var livestreams = await apiClient.GetScheduledLivestreamsAsync(businessUnit, "SPORT", cancellationToken).ConfigureAwait(false);
if (livestreams != null)
{
all.AddRange(livestreams);
}
}
// Filter to only future/current livestreams that aren't blocked
return all
.Where(ls => ls.Blocked != true && (ls.ValidTo == null || ls.ValidTo.Value.ToUniversalTime() > DateTime.UtcNow))
.OrderBy(ls => ls.ValidFrom)
.ToList();
}
///
public async Task ScheduleRecordingAsync(string urn, CancellationToken cancellationToken)
{
await LoadRecordingsAsync().ConfigureAwait(false);
// Check if already scheduled
var existing = _recordings.FirstOrDefault(r => r.Urn == urn && r.State is RecordingState.Scheduled or RecordingState.WaitingForStream or RecordingState.Recording);
if (existing != null)
{
_logger.LogInformation("Recording already exists for URN {Urn} in state {State}", urn, existing.State);
return existing;
}
// Fetch metadata for the URN. The unit is encoded in the URN itself, so use that to
// query the right schedule regardless of which units are enabled.
var config = Plugin.Instance?.Configuration;
var businessUnit = ParseBusinessUnitFromUrn(urn)
?? (config?.BusinessUnit ?? Configuration.BusinessUnit.SRF).ToString().ToLowerInvariant();
using var apiClient = _apiClientFactory.CreateClient();
var livestreams = await apiClient.GetScheduledLivestreamsAsync(businessUnit, "SPORT", cancellationToken).ConfigureAwait(false);
var program = livestreams?.FirstOrDefault(ls => ls.Urn == urn);
var entry = new RecordingEntry
{
Id = Guid.NewGuid().ToString("N"),
Urn = urn,
BusinessUnit = ParseBusinessUnitFromUrn(urn) ?? businessUnit,
Title = program?.Title ?? urn,
Description = program?.Lead ?? program?.Description,
ImageUrl = program?.ImageUrl,
ValidFrom = program?.ValidFrom,
ValidTo = program?.ValidTo,
State = RecordingState.Scheduled,
CreatedAt = DateTime.UtcNow
};
_recordings.Add(entry);
await SaveRecordingsAsync().ConfigureAwait(false);
_logger.LogInformation("Scheduled recording for '{Title}' (URN: {Urn}, starts: {ValidFrom})", entry.Title, urn, entry.ValidFrom);
return entry;
}
///
/// Extracts the business unit from an SRF URN of the form "urn:<bu>:<type>:...".
///
/// The URN.
/// The lowercase business unit, or null if it cannot be determined.
private static string? ParseBusinessUnitFromUrn(string urn)
{
if (string.IsNullOrEmpty(urn))
{
return null;
}
var parts = urn.Split(':');
return parts.Length >= 2 && !string.IsNullOrEmpty(parts[1])
? parts[1].ToLowerInvariant()
: null;
}
///
public bool CancelRecording(string recordingId)
{
var entry = _recordings.FirstOrDefault(r => r.Id == recordingId);
if (entry == null)
{
return false;
}
if (entry.State == RecordingState.Recording)
{
StopFfmpeg(recordingId);
}
entry.State = RecordingState.Cancelled;
entry.RecordingEndedAt = DateTime.UtcNow;
_ = SaveRecordingsAsync();
_logger.LogInformation("Cancelled recording '{Title}' ({Id})", entry.Title, recordingId);
return true;
}
///
public bool StopRecording(string recordingId)
{
var entry = _recordings.FirstOrDefault(r => r.Id == recordingId && r.State == RecordingState.Recording);
if (entry == null)
{
return false;
}
StopFfmpeg(recordingId);
CompleteRecording(entry, DateTime.UtcNow);
_ = SaveRecordingsAsync();
_logger.LogInformation("Stopped recording '{Title}' ({Id})", entry.Title, recordingId);
return true;
}
///
public IReadOnlyList GetRecordings(RecordingState? stateFilter)
{
// Ensure loaded synchronously for simple reads
if (!_loaded)
{
LoadRecordingsAsync().GetAwaiter().GetResult();
}
if (stateFilter.HasValue)
{
return _recordings.Where(r => r.State == stateFilter.Value).OrderByDescending(r => r.CreatedAt).ToList();
}
return _recordings.OrderByDescending(r => r.CreatedAt).ToList();
}
///
public RecordingEntry? GetRecording(string recordingId)
{
return GetRecordings(null).FirstOrDefault(r => r.Id == recordingId);
}
///
public IReadOnlyList 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();
}
///
public bool DeleteRecording(string recordingId, bool deleteFile)
{
var entry = _recordings.FirstOrDefault(r => r.Id == recordingId);
if (entry == null)
{
return false;
}
if (entry.State == RecordingState.Recording)
{
StopFfmpeg(recordingId);
}
if (deleteFile && !string.IsNullOrEmpty(entry.OutputPath) && File.Exists(entry.OutputPath))
{
try
{
File.Delete(entry.OutputPath);
_logger.LogInformation("Deleted recording file: {Path}", entry.OutputPath);
}
catch (Exception ex)
{
_logger.LogWarning(ex, "Failed to delete recording file: {Path}", entry.OutputPath);
}
}
_recordings.Remove(entry);
_ = SaveRecordingsAsync();
_logger.LogInformation("Deleted recording entry '{Title}' ({Id})", entry.Title, recordingId);
return true;
}
///
public async Task ProcessRecordingsAsync(CancellationToken cancellationToken)
{
// Prevent overlapping scheduler runs from spawning duplicate ffmpeg processes
if (!await _processLock.WaitAsync(0, CancellationToken.None).ConfigureAwait(false))
{
_logger.LogDebug("ProcessRecordingsAsync already running, skipping");
return;
}
try
{
await ProcessRecordingsCoreAsync(cancellationToken).ConfigureAwait(false);
}
finally
{
_processLock.Release();
}
}
private async Task ProcessRecordingsCoreAsync(CancellationToken cancellationToken)
{
await LoadRecordingsAsync().ConfigureAwait(false);
var now = DateTime.UtcNow;
var changed = false;
foreach (var entry in _recordings.ToList())
{
// Normalize ValidFrom/ValidTo to UTC for correct comparison
var validFromUtc = entry.ValidFrom.HasValue ? entry.ValidFrom.Value.ToUniversalTime() : (DateTime?)null;
var validToUtc = entry.ValidTo.HasValue ? entry.ValidTo.Value.ToUniversalTime() : (DateTime?)null;
switch (entry.State)
{
case RecordingState.Scheduled:
case RecordingState.WaitingForStream:
// Give up once the broadcast window is over, otherwise we'd retry forever
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(
"Time to start recording '{Title}': ValidFrom={ValidFrom} (UTC: {ValidFromUtc}), Now={Now}",
entry.Title,
entry.ValidFrom,
validFromUtc,
now);
changed |= await TryStartRecordingAsync(entry, cancellationToken).ConfigureAwait(false);
}
break;
case RecordingState.Recording:
// Check if recording should stop (ValidTo reached or process died)
if (validToUtc.HasValue && validToUtc.Value <= now)
{
_logger.LogInformation("Recording '{Title}' reached ValidTo, stopping", entry.Title);
StopFfmpeg(entry.Id);
CompleteRecording(entry, now);
changed = true;
}
else if (!_activeProcesses.ContainsKey(entry.Id))
{
// ffmpeg process died unexpectedly — try to restart
_logger.LogWarning("ffmpeg process for '{Title}' is no longer running, attempting restart", entry.Title);
changed |= await TryStartRecordingAsync(entry, cancellationToken).ConfigureAwait(false);
}
break;
}
}
if (changed)
{
await SaveRecordingsAsync().ConfigureAwait(false);
}
}
private async Task TryStartRecordingAsync(RecordingEntry entry, CancellationToken cancellationToken)
{
try
{
// Fetch the media composition to get the stream URL
var mediaComposition = await _mediaCompositionFetcher.GetMediaCompositionAsync(entry.Urn, cacheDurationOverride: 2, cancellationToken: cancellationToken).ConfigureAwait(false);
var chapter = mediaComposition?.ChapterList is { Count: > 0 } list ? list[0] : null;
if (chapter == null)
{
_logger.LogDebug("No chapter found for '{Title}', stream may not be live yet", entry.Title);
MarkStreamUnavailable(entry);
return true;
}
var config = Plugin.Instance?.Configuration;
var quality = config?.QualityPreference ?? Configuration.QualityPreference.Auto;
var streamUrl = _streamUrlResolver.GetStreamUrl(chapter, quality);
if (string.IsNullOrEmpty(streamUrl))
{
_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;
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
var itemId = $"rec_{entry.Id}";
var isLiveStream = chapter.Type == "SCHEDULED_LIVESTREAM" || UrnHelper.IsLivestreamUrn(entry.Urn);
_proxyService.RegisterStreamDeferred(itemId, streamUrl, entry.Urn, isLiveStream);
// Build proxy URL for ffmpeg (use localhost for local access)
var proxyUrl = $"{GetServerBaseUrl()}/Plugins/SRFPlay/Proxy/{itemId}/master.m3u8";
// 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 timestamp = DisplayTime.Now().ToString("yyyy-MM-dd_HHmm", CultureInfo.InvariantCulture);
outputPath = Path.Combine(GetRecordingOutputPath(), $"{safeTitle}_{timestamp}.ts");
entry.OutputPath = outputPath;
}
// Start ffmpeg
StartFfmpeg(entry.Id, proxyUrl, outputPath);
entry.State = RecordingState.Recording;
entry.RecordingStartedAt ??= DateTime.UtcNow;
_logger.LogInformation("Started recording '{Title}' to {OutputPath}", entry.Title, outputPath);
return true;
}
catch (Exception ex) when (ex is not OperationCanceledException)
{
_logger.LogError(ex, "Failed to start recording '{Title}'", entry.Title);
entry.State = RecordingState.Failed;
entry.ErrorMessage = ex.Message;
return true;
}
}
///
/// 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).
///
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;
}
}
///
/// Reads the DVR window start (unix seconds in the start= query parameter) from a stream URL.
///
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)
{
var process = new Process
{
StartInfo = new ProcessStartInfo
{
FileName = _mediaEncoder.EncoderPath,
// 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,
RedirectStandardInput = true,
RedirectStandardOutput = true,
RedirectStandardError = true,
CreateNoWindow = true
}
};
process.ErrorDataReceived += (_, args) =>
{
if (!string.IsNullOrEmpty(args.Data))
{
_logger.LogDebug("ffmpeg [{RecordingId}]: {Data}", recordingId, args.Data);
}
};
process.Start();
process.BeginErrorReadLine();
var active = new ActiveRecording(process);
_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);
}
///
/// 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.
///
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(recordingId, active));
}
}
private void StopFfmpeg(string recordingId)
{
if (_activeProcesses.TryRemove(recordingId, out var active))
{
var process = active.Process;
try
{
if (!process.HasExited)
{
// Send 'q' to ffmpeg stdin for graceful shutdown
process.StandardInput.Write("q");
process.StandardInput.Flush();
if (!process.WaitForExit(10000))
{
_logger.LogWarning("ffmpeg did not exit gracefully for {RecordingId}, killing", recordingId);
process.Kill(true);
}
}
// Let the remaining output reach the file before releasing the process
active.Pump?.Wait(TimeSpan.FromSeconds(10));
process.Dispose();
}
catch (Exception ex)
{
_logger.LogWarning(ex, "Error stopping ffmpeg for recording {RecordingId}", recordingId);
}
}
}
private static string SanitizeFileName(string name)
{
var invalid = Path.GetInvalidFileNameChars();
var sanitized = string.Join("_", name.Split(invalid, StringSplitOptions.RemoveEmptyEntries));
// Also replace spaces and other problematic chars
sanitized = Regex.Replace(sanitized, @"[\s]+", "_");
return sanitized.Length > 100 ? sanitized[..100] : sanitized;
}
///
public void Dispose()
{
Dispose(true);
GC.SuppressFinalize(this);
}
///
/// Releases resources.
///
/// True to release managed resources.
protected virtual void Dispose(bool disposing)
{
if (_disposed)
{
return;
}
if (disposing)
{
foreach (var kvp in _activeProcesses)
{
StopFfmpeg(kvp.Key);
}
_persistLock.Dispose();
_processLock.Dispose();
}
_disposed = true;
}
///
/// A running ffmpeg process and the task copying its output to disk.
///
private sealed class ActiveRecording
{
public ActiveRecording(Process process)
{
Process = process;
}
public Process Process { get; }
public Task? Pump { get; set; }
}
}