Files
jellyfin-srfPlay/Jellyfin.Plugin.SRFPlay/Services/RecordingService.cs
T
dtourolleandClaude Opus 5.5 aa47de9912
🏗️ Build Plugin / build (push) Successful in 12m31s
Nightly Build / nightly-build (push) Successful in 2m44s
🧪 Test Plugin / test (push) Successful in 11m14s
🚀 Release Plugin / build-and-release (push) Successful in 12m33s
Record into a Jellyfin library by default; configurable recording location
When no custom directory is set, recordings go to an "SRF Recordings" folder
in the chosen library, or the first TV Shows library, so they show up without
adding a library by hand. The settings page gets a library picker, a check
that the directory is writable, and help for sandboxed systemd services.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-26 19:35:55 -04:00

810 lines
30 KiB
C#

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.Library;
using MediaBrowser.Controller.MediaEncoding;
using Microsoft.Extensions.Logging;
namespace Jellyfin.Plugin.SRFPlay.Services;
/// <summary>
/// Service for managing sport livestream recordings using ffmpeg.
/// </summary>
public class RecordingService : IRecordingService, IDisposable
{
/// <summary>Folder created inside a library for recordings; shows up there as one series.</summary>
private const string RecordingsFolderName = "SRF Recordings";
/// <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 ISRFApiClientFactory _apiClientFactory;
private readonly IStreamProxyService _proxyService;
private readonly IStreamUrlResolver _streamUrlResolver;
private readonly IMediaCompositionFetcher _mediaCompositionFetcher;
private readonly IServerApplicationHost _appHost;
private readonly IMediaEncoder _mediaEncoder;
private readonly ILibraryManager _libraryManager;
private readonly ConcurrentDictionary<string, ActiveRecording> _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<RecordingEntry> _recordings = new();
private bool _loaded;
private bool _disposed;
/// <summary>
/// Initializes a new instance of the <see cref="RecordingService"/> class.
/// </summary>
/// <param name="logger">The logger.</param>
/// <param name="apiClientFactory">The API client factory.</param>
/// <param name="proxyService">The stream proxy service.</param>
/// <param name="streamUrlResolver">The stream URL resolver.</param>
/// <param name="mediaCompositionFetcher">The media composition fetcher.</param>
/// <param name="appHost">The application host.</param>
/// <param name="mediaEncoder">The media encoder for ffmpeg path.</param>
/// <param name="libraryManager">The library manager, used to record into a library.</param>
public RecordingService(
ILogger<RecordingService> logger,
ISRFApiClientFactory apiClientFactory,
IStreamProxyService proxyService,
IStreamUrlResolver streamUrlResolver,
IMediaCompositionFetcher mediaCompositionFetcher,
IServerApplicationHost appHost,
IMediaEncoder mediaEncoder,
ILibraryManager libraryManager)
{
_logger = logger;
_apiClientFactory = apiClientFactory;
_proxyService = proxyService;
_streamUrlResolver = streamUrlResolver;
_mediaCompositionFetcher = mediaCompositionFetcher;
_appHost = appHost;
_mediaEncoder = mediaEncoder;
_libraryManager = libraryManager;
}
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 path = ResolveOutputDirectory().Path;
Directory.CreateDirectory(path);
return path;
}
/// <summary>
/// Picks the recording directory: custom path, then the chosen library, then the first
/// TV Shows library, then a folder in the service user's home.
/// </summary>
private RecordingOutputInfo ResolveOutputDirectory()
{
var config = Plugin.Instance?.Configuration;
if (!string.IsNullOrWhiteSpace(config?.RecordingOutputPath))
{
return new RecordingOutputInfo { Path = config.RecordingOutputPath.Trim(), Source = "custom" };
}
try
{
var libraries = _libraryManager.GetVirtualFolders()
.Where(f => f.Locations != null && f.Locations.Length > 0)
.ToList();
var selectedId = config?.RecordingLibraryId;
var library = string.IsNullOrWhiteSpace(selectedId)
? null
: libraries.FirstOrDefault(f => string.Equals(f.ItemId, selectedId, StringComparison.OrdinalIgnoreCase));
var source = "library";
if (library == null)
{
// CollectionType is an enum whose member casing differs between Jellyfin versions
library = libraries.FirstOrDefault(f => string.Equals(f.CollectionType?.ToString(), "tvshows", StringComparison.OrdinalIgnoreCase));
source = "tvLibrary";
}
if (library != null)
{
return new RecordingOutputInfo
{
Path = Path.Combine(library.Locations[0], RecordingsFolderName),
Source = source,
LibraryName = library.Name
};
}
}
catch (Exception ex)
{
_logger.LogWarning(ex, "Could not read libraries to pick a recording directory");
}
return new RecordingOutputInfo
{
Path = Path.Combine(Environment.GetFolderPath(Environment.SpecialFolder.UserProfile), RecordingsFolderName),
Source = "fallback"
};
}
/// <inheritdoc />
public RecordingOutputInfo GetOutputInfo()
{
var info = ResolveOutputDirectory();
try
{
Directory.CreateDirectory(info.Path);
var probe = Path.Combine(info.Path, $".srfplay-write-test-{Guid.NewGuid():N}");
File.WriteAllText(probe, string.Empty);
File.Delete(probe);
info.Writable = true;
}
catch (Exception ex) when (ex is IOException or UnauthorizedAccessException)
{
info.Error = ex.Message;
}
return info;
}
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<List<RecordingEntry>>(json) ?? new List<RecordingEntry>();
_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<RecordingEntry>();
}
}
_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();
}
}
/// <inheritdoc />
public async Task<IReadOnlyList<PlayV3TvProgram>> GetUpcomingScheduleAsync(CancellationToken cancellationToken)
{
var units = System.Enum.GetValues<Configuration.BusinessUnit>();
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<PlayV3TvProgram>();
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();
}
/// <inheritdoc />
public async Task<RecordingEntry> 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;
}
/// <summary>
/// Extracts the business unit from an SRF URN of the form "urn:&lt;bu&gt;:&lt;type&gt;:...".
/// </summary>
/// <param name="urn">The URN.</param>
/// <returns>The lowercase business unit, or null if it cannot be determined.</returns>
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;
}
/// <inheritdoc />
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;
}
/// <inheritdoc />
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;
}
/// <inheritdoc />
public IReadOnlyList<RecordingEntry> 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();
}
/// <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 />
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;
}
/// <inheritdoc />
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<bool> 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;
}
}
/// <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)
{
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);
}
/// <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)
{
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;
}
/// <inheritdoc />
public void Dispose()
{
Dispose(true);
GC.SuppressFinalize(this);
}
/// <summary>
/// Releases resources.
/// </summary>
/// <param name="disposing">True to release managed resources.</param>
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;
}
/// <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; }
}
}