TsDemuxer.cs 8.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212
  1. using System;
  2. using System.Collections.Generic;
  3. using System.Collections.ObjectModel;
  4. using System.Diagnostics;
  5. using System.Linq;
  6. using System.Text;
  7. using MatrixIO.IO.MpegTs.Streams;
  8. using MatrixIO.IO.MpegTs.Tables;
  9. namespace MatrixIO.IO.MpegTs
  10. {
  11. public class TsDemuxer
  12. {
  13. private readonly object _syncObject = new object();
  14. public IList<TsProgram> Programs { get; private set; }
  15. private readonly Dictionary<ushort, TsStream> _streams = new Dictionary<ushort, TsStream>();
  16. public ICollection<TsStream> Streams
  17. {
  18. get { return _streams.Values; }
  19. }
  20. private int _partialPacketLength = 0;
  21. private readonly byte[] _partialPacket = new byte[TsPacket.Length];
  22. public TsDemuxer()
  23. {
  24. Programs = Portability.CreateList<TsProgram>();
  25. var programAssociationStream = new TableStream()
  26. {PacketIdentifier = (ushort) PacketIdentifier.ProgramAssociationTable};
  27. programAssociationStream.UnitReceived += ProgramAssociationTableReceived;
  28. _streams.Add((ushort) PacketIdentifier.ProgramAssociationTable, programAssociationStream);
  29. var descriptionStream = new TableStream() {PacketIdentifier = (ushort) PacketIdentifier.TsDescriptionTable};
  30. descriptionStream.UnitReceived += DescriptionTableReceived;
  31. _streams.Add((ushort) PacketIdentifier.TsDescriptionTable, descriptionStream);
  32. }
  33. public void ProcessInput(byte[] buffer, int offset = 0)
  34. {
  35. ProcessInput(buffer, offset, buffer.Length);
  36. }
  37. public void ProcessInput(byte[] buffer, int offset, int length)
  38. {
  39. int position = offset;
  40. int remainder = 0;
  41. if (_partialPacketLength != 0)
  42. {
  43. Debug.WriteLine("Using previous " + _partialPacketLength + " byte partial packet.");
  44. remainder = TsPacket.Length - _partialPacketLength;
  45. int len = Math.Min(remainder, length);
  46. Debug.WriteLine("Copying " + len + " additional bytes to partial packet.");
  47. Buffer.BlockCopy(buffer, 0, _partialPacket, _partialPacketLength, len);
  48. _partialPacketLength += len;
  49. position += len;
  50. Debug.WriteLine("Partial packet is now " + _partialPacketLength + " bytes.");
  51. if (_partialPacketLength >= TsPacket.Length)
  52. {
  53. _partialPacketLength = 0;
  54. try
  55. {
  56. var packet = new TsPacket(_partialPacket);
  57. ProcessPacket(packet);
  58. Debug.WriteLine("Reassembled partial packet.");
  59. }
  60. catch (Exception ex)
  61. {
  62. Debug.WriteLine(ex.ToString());
  63. }
  64. }
  65. }
  66. while ((offset + length - position) >= TsPacket.Length)
  67. {
  68. if (buffer[position] == 0x47)
  69. {
  70. try
  71. {
  72. var packet = new TsPacket(buffer, position);
  73. position += TsPacket.Length;
  74. Debug.WriteLine("Processing packet.");
  75. ProcessPacket(packet);
  76. }
  77. catch
  78. {
  79. Debug.WriteLine("Error deserializing packet. Assuming false start code and scanning ahead for the next packet.");
  80. position++;
  81. }
  82. }
  83. else
  84. {
  85. Debug.WriteLine("Skipping byte 0x" + buffer[position].ToString("X2") + " at offset " + position + ".");
  86. position++;
  87. }
  88. }
  89. remainder = offset + length - position;
  90. if (remainder > 0)
  91. {
  92. Debug.WriteLine("Remainder is " + remainder + " bytes.");
  93. int packetStartOffset = -1;
  94. for (int i = 0; i < remainder; i++)
  95. {
  96. if (buffer[position + i] == 0x47)
  97. {
  98. packetStartOffset = i;
  99. break;
  100. }
  101. Debug.WriteLine("Skipping byte 0x" + buffer[position].ToString("X2") + " at offset " + (position + i) + ".");
  102. }
  103. if (packetStartOffset >= 0)
  104. {
  105. Debug.WriteLine("Found sync byte in remainder. Storing partial packet.");
  106. _partialPacketLength = remainder - packetStartOffset;
  107. Buffer.BlockCopy(buffer, position + packetStartOffset, _partialPacket, 0, _partialPacketLength);
  108. }
  109. else _partialPacketLength = 0;
  110. }
  111. }
  112. public void ProcessPacket(TsPacket packet)
  113. {
  114. if (packet.PacketIdentifier == (ushort) PacketIdentifier.NullPacket)
  115. return; // null packet for padding strict muxrate streams
  116. #if DEBUG2
  117. var dbg = new StringBuilder();
  118. dbg.Append("Received packet for PID " + (PacketIdentifier) packet.PacketIdentifier + ". ");
  119. dbg.Append("Continuity is " + packet.ContinuityCounter + ". ");
  120. if (packet.AdaptationField != null) dbg.Append("Has Adaptation Field. ");
  121. if (packet.Payload != null) dbg.Append("Has " + packet.Payload.Length + " byte Payload. ");
  122. if (packet.PayloadUnitStartIndicator) dbg.Append("First in series. ");
  123. Debug.WriteLine(dbg.ToString());
  124. #endif
  125. TsStream stream;
  126. if (_streams.TryGetValue(packet.PacketIdentifier, out stream)) stream.ProcessInput(packet);
  127. }
  128. private void ProgramAssociationTableReceived(object sender, TsStreamEventArgs eventArgs)
  129. {
  130. var e = eventArgs as TsStreamEventArgs<TsTable>;
  131. if (e == null) return;
  132. lock (_syncObject)
  133. {
  134. var pat = e.DecodedUnit.TableIdentifier == TableIdentifier.ProgramAssociation
  135. ? (ProgramAssociationTable) e.DecodedUnit
  136. : null;
  137. if (pat == null) return;
  138. #if DEBUG2
  139. Debug.WriteLine("Table is a " + pat.TableIdentifier + " with Identifier 0x" + pat.Identifier.ToString("X4") +
  140. " and IsCurrent=" + pat.IsCurrent);
  141. foreach (var pa in pat.Rows)
  142. Debug.WriteLine("Program " + pa.ProgramNumber + " is on PID " + pa.PacketIdentifier);
  143. #endif
  144. var updatedPrograms = (from r in pat.Rows select r.ProgramNumber).ToArray();
  145. var existingPrograms = (from p in Programs select p.ProgramNumber).ToArray();
  146. var newProgramRows = from r in pat.Rows
  147. where updatedPrograms.Except(existingPrograms).Contains(r.ProgramNumber)
  148. select r;
  149. foreach (var r in newProgramRows)
  150. {
  151. var newProgram = new TsProgram(r.ProgramNumber,
  152. new TableStream() {PacketIdentifier = r.PacketIdentifier})
  153. {
  154. Status =
  155. pat.IsCurrent ? TsProgramStatus.Current : TsProgramStatus.Next
  156. };
  157. _streams.Add(r.PacketIdentifier, newProgram.ProgramMapStream);
  158. Programs.Add(newProgram);
  159. newProgram.StreamAdded += AddProgramStream;
  160. }
  161. if (pat.IsCurrent)
  162. {
  163. var oldPrograms = from p in Programs
  164. where !updatedPrograms.Contains(p.ProgramNumber)
  165. select p;
  166. foreach (var p in oldPrograms)
  167. {
  168. p.Status = TsProgramStatus.Dicontinued;
  169. }
  170. }
  171. }
  172. }
  173. private void DescriptionTableReceived(object sender, TsStreamEventArgs eventArgs)
  174. {
  175. var e = eventArgs as TsStreamEventArgs<TsTable>;
  176. if (e == null) return;
  177. Debug.WriteLine("Received Description Table");
  178. }
  179. private void AddProgramStream(object sender, ProgramStreamEventArgs e)
  180. {
  181. _streams.Add(e.Stream.PacketIdentifier, e.Stream);
  182. }
  183. }
  184. }