aboutsummaryrefslogtreecommitdiff
path: root/Jellyfin.Api/Helpers/ProgressiveFileCopier.cs
blob: 8bddf00d5f2b584e596f335a3b6df30777a6b71a (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
using System;
using System.Buffers;
using System.IO;
using System.Runtime.InteropServices;
using System.Threading;
using System.Threading.Tasks;
using Jellyfin.Api.Models.PlaybackDtos;
using MediaBrowser.Common.Extensions;
using MediaBrowser.Controller.Library;
using MediaBrowser.Model.IO;

namespace Jellyfin.Api.Helpers
{
    /// <summary>
    /// Progressive file copier.
    /// </summary>
    public class ProgressiveFileCopier
    {
        private readonly TranscodingJobDto? _job;
        private readonly string? _path;
        private readonly CancellationToken _cancellationToken;
        private readonly IDirectStreamProvider? _directStreamProvider;
        private readonly TranscodingJobHelper _transcodingJobHelper;
        private long _bytesWritten;

        /// <summary>
        /// Initializes a new instance of the <see cref="ProgressiveFileCopier"/> class.
        /// </summary>
        /// <param name="path">The path to copy from.</param>
        /// <param name="job">The transcoding job.</param>
        /// <param name="transcodingJobHelper">Instance of the <see cref="TranscodingJobHelper"/>.</param>
        /// <param name="cancellationToken">The cancellation token.</param>
        public ProgressiveFileCopier(string path, TranscodingJobDto? job, TranscodingJobHelper transcodingJobHelper, CancellationToken cancellationToken)
        {
            _path = path;
            _job = job;
            _cancellationToken = cancellationToken;
            _transcodingJobHelper = transcodingJobHelper;
        }

        /// <summary>
        /// Initializes a new instance of the <see cref="ProgressiveFileCopier"/> class.
        /// </summary>
        /// <param name="directStreamProvider">Instance of the <see cref="IDirectStreamProvider"/> interface.</param>
        /// <param name="job">The transcoding job.</param>
        /// <param name="transcodingJobHelper">Instance of the <see cref="TranscodingJobHelper"/>.</param>
        /// <param name="cancellationToken">The cancellation token.</param>
        public ProgressiveFileCopier(IDirectStreamProvider directStreamProvider, TranscodingJobDto? job, TranscodingJobHelper transcodingJobHelper, CancellationToken cancellationToken)
        {
            _directStreamProvider = directStreamProvider;
            _job = job;
            _cancellationToken = cancellationToken;
            _transcodingJobHelper = transcodingJobHelper;
        }

        /// <summary>
        /// Gets or sets a value indicating whether allow read end of file.
        /// </summary>
        public bool AllowEndOfFile { get; set; } = true;

        /// <summary>
        /// Gets or sets copy start position.
        /// </summary>
        public long StartPosition { get; set; }

        /// <summary>
        /// Write source stream to output.
        /// </summary>
        /// <param name="outputStream">Output stream.</param>
        /// <param name="cancellationToken">Cancellation token.</param>
        /// <returns>A <see cref="Task"/>.</returns>
        public async Task WriteToAsync(Stream outputStream, CancellationToken cancellationToken)
        {
            cancellationToken = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, _cancellationToken).Token;

            try
            {
                if (_directStreamProvider != null)
                {
                    await _directStreamProvider.CopyToAsync(outputStream, cancellationToken).ConfigureAwait(false);
                    return;
                }

                var fileOptions = FileOptions.SequentialScan;
                var allowAsyncFileRead = false;

                // use non-async filestream along with read due to https://github.com/dotnet/corefx/issues/6039
                if (!RuntimeInformation.IsOSPlatform(OSPlatform.Windows))
                {
                    fileOptions |= FileOptions.Asynchronous;
                    allowAsyncFileRead = true;
                }

                if (_path == null)
                {
                    throw new ResourceNotFoundException(nameof(_path));
                }

                await using var inputStream = new FileStream(_path, FileMode.Open, FileAccess.Read, FileShare.ReadWrite, IODefaults.FileStreamBufferSize, fileOptions);

                var eofCount = 0;
                const int EmptyReadLimit = 20;
                if (StartPosition > 0)
                {
                    inputStream.Position = StartPosition;
                }

                while (eofCount < EmptyReadLimit || !AllowEndOfFile)
                {
                    var bytesRead = await CopyToInternalAsync(inputStream, outputStream, allowAsyncFileRead, cancellationToken).ConfigureAwait(false);

                    if (bytesRead == 0)
                    {
                        if (_job == null || _job.HasExited)
                        {
                            eofCount++;
                        }

                        await Task.Delay(100, cancellationToken).ConfigureAwait(false);
                    }
                    else
                    {
                        eofCount = 0;
                    }
                }
            }
            finally
            {
                if (_job != null)
                {
                    _transcodingJobHelper.OnTranscodeEndRequest(_job);
                }
            }
        }

        private async Task<int> CopyToInternalAsync(Stream source, Stream destination, bool readAsync, CancellationToken cancellationToken)
        {
            var array = ArrayPool<byte>.Shared.Rent(IODefaults.CopyToBufferSize);
            try
            {
                int bytesRead;
                int totalBytesRead = 0;

                if (readAsync)
                {
                    bytesRead = await source.ReadAsync(array, 0, array.Length, cancellationToken).ConfigureAwait(false);
                }
                else
                {
                    bytesRead = source.Read(array, 0, array.Length);
                }

                while (bytesRead != 0)
                {
                    var bytesToWrite = bytesRead;

                    if (bytesToWrite > 0)
                    {
                        await destination.WriteAsync(array, 0, Convert.ToInt32(bytesToWrite), cancellationToken).ConfigureAwait(false);

                        _bytesWritten += bytesRead;
                        totalBytesRead += bytesRead;

                        if (_job != null)
                        {
                            _job.BytesDownloaded = Math.Max(_job.BytesDownloaded ?? _bytesWritten, _bytesWritten);
                        }
                    }

                    if (readAsync)
                    {
                        bytesRead = await source.ReadAsync(array, 0, array.Length, cancellationToken).ConfigureAwait(false);
                    }
                    else
                    {
                        bytesRead = source.Read(array, 0, array.Length);
                    }
                }

                return totalBytesRead;
            }
            finally
            {
                ArrayPool<byte>.Shared.Return(array);
            }
        }
    }
}