QueueStream.cs 2.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101
  1. using System;
  2. using System.Collections.Concurrent;
  3. using System.Collections.Generic;
  4. using System.IO;
  5. using System.Linq;
  6. using System.Text;
  7. using System.Threading;
  8. using System.Threading.Tasks;
  9. using MediaBrowser.Model.Logging;
  10. namespace Emby.Server.Implementations.LiveTv.TunerHosts
  11. {
  12. public class QueueStream
  13. {
  14. private readonly Stream _outputStream;
  15. private readonly ConcurrentQueue<byte[]> _queue = new ConcurrentQueue<byte[]>();
  16. private CancellationToken _cancellationToken;
  17. public TaskCompletionSource<bool> TaskCompletion { get; private set; }
  18. public Action<QueueStream> OnFinished { get; set; }
  19. private readonly ILogger _logger;
  20. private bool _isActive;
  21. public QueueStream(Stream outputStream, ILogger logger)
  22. {
  23. _outputStream = outputStream;
  24. _logger = logger;
  25. TaskCompletion = new TaskCompletionSource<bool>();
  26. }
  27. public void Queue(byte[] bytes)
  28. {
  29. if (_isActive)
  30. {
  31. _queue.Enqueue(bytes);
  32. }
  33. }
  34. public void Start(CancellationToken cancellationToken)
  35. {
  36. _cancellationToken = cancellationToken;
  37. Task.Run(() => StartInternal());
  38. }
  39. private byte[] Dequeue()
  40. {
  41. byte[] bytes;
  42. if (_queue.TryDequeue(out bytes))
  43. {
  44. return bytes;
  45. }
  46. return null;
  47. }
  48. private async Task StartInternal()
  49. {
  50. var cancellationToken = _cancellationToken;
  51. try
  52. {
  53. while (!cancellationToken.IsCancellationRequested)
  54. {
  55. _isActive = true;
  56. var bytes = Dequeue();
  57. if (bytes != null)
  58. {
  59. await _outputStream.WriteAsync(bytes, 0, bytes.Length, cancellationToken).ConfigureAwait(false);
  60. }
  61. else
  62. {
  63. await Task.Delay(50, cancellationToken).ConfigureAwait(false);
  64. }
  65. }
  66. TaskCompletion.TrySetResult(true);
  67. _logger.Debug("QueueStream complete");
  68. }
  69. catch (OperationCanceledException)
  70. {
  71. _logger.Debug("QueueStream cancelled");
  72. TaskCompletion.TrySetCanceled();
  73. }
  74. catch (Exception ex)
  75. {
  76. _logger.ErrorException("Error in QueueStream", ex);
  77. TaskCompletion.TrySetException(ex);
  78. }
  79. finally
  80. {
  81. _isActive = false;
  82. if (OnFinished != null)
  83. {
  84. OnFinished(this);
  85. }
  86. }
  87. }
  88. }
  89. }