reader.js 6.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224
  1. // stream from hls source
  2. var util = require('util'),
  3. url = require('url'),
  4. zlib = require('zlib'),
  5. assert = require('assert');
  6. var request = require('request'),
  7. oncemore = require('./oncemore'),
  8. uristream = require('./uristream'),
  9. debug = require('debug')('hls:reader');
  10. try {
  11. var Readable = require('stream').Readable;
  12. assert(Readable);
  13. } catch (e) {
  14. var Readable = require('readable-stream');
  15. }
  16. var m3u8 = require('./m3u8');
  17. function noop() {};
  18. module.exports = hlsreader;
  19. hlsreader.HlsStreamReader = HlsStreamReader;
  20. /*
  21. options:
  22. startSeq*
  23. noData // don't emit any data - useful for analyzing the stream structure
  24. maxRedirects*
  25. cacheDir*
  26. headers* // allows for custom user-agent, cookies, auth, etc
  27. emits:
  28. index (m3u8)
  29. segment (seqNo, duration, datetime, size?, )
  30. */
  31. function HlsStreamReader(src, options) {
  32. var self = this;
  33. options = options || {};
  34. if (typeof src === 'string')
  35. src = url.parse(src);
  36. this.url = src;
  37. this.baseUrl = src;
  38. this.fullStream = !!options.fullStream;
  39. this.keepConnection = !!options.keepConnection;
  40. this.noData = !!options.noData;
  41. this.indexStream = null;
  42. this.index = null;
  43. this.readState = {
  44. currentSeq:-1,
  45. currentSegment:null,
  46. stream:null
  47. }
  48. function getUpdateInterval(updated) {
  49. if (updated && self.index.segments.length)
  50. return Math.min(self.index.target_duration, self.index.segments[self.index.segments.length-1].duration);
  51. else
  52. return self.index.target_duration / 2;
  53. }
  54. function updatecheck(updated) {
  55. if (updated) {
  56. if (self.readState.currentSeq===-1)
  57. self.readState.currentSeq = self.index.startSeqNo(self.fullStream);
  58. else if (self.readState.currentSeq < self.index.startSeqNo(true))
  59. self.readState.currentSeq = self.index.startSeqNo(true);
  60. self.emit('index', self.index);
  61. if (self.index.variant)
  62. return self.end();
  63. }
  64. checkcurrent();
  65. if (!self.index.ended) {
  66. var updateInterval = getUpdateInterval(updated);
  67. debug('scheduling index refresh', updateInterval);
  68. setTimeout(updateindex, Math.max(1, updateInterval)*1000);
  69. }
  70. }
  71. function updateindex() {
  72. var stream = uristream(url.format(self.url), { timeout:30*1000 });
  73. stream.on('meta', function(meta) {
  74. if (meta.mime !== 'application/vnd.apple.mpegurl' &&
  75. meta.mime !== 'application/x-mpegurl' && meta.mime !== 'audio/mpegurl')
  76. return stream.emit('error', new Error('Invalid MIME type: '+meta.mime));
  77. // FIXME: correctly handle .m3u us-ascii encoding
  78. self.baseUrl = meta.url;
  79. });
  80. m3u8.parse(stream, function(err, index) {
  81. if (err) {
  82. if (self.index && self.keepConnection) {
  83. console.error('Failed to parse index at '+url.format(self.url)+':', err.stack || err);
  84. return updatecheck(false);
  85. }
  86. return self.emit('error', err);
  87. }
  88. var updated = true;
  89. if (self.index && self.index.lastSeqNo() === index.lastSeqNo()) {
  90. debug('index was not updated');
  91. updated = false;
  92. }
  93. self.index = index;
  94. updatecheck(updated);
  95. });
  96. }
  97. updateindex();
  98. function checkcurrent() {
  99. if (self.readState.currentSegment) return; // already processing
  100. self.readState.currentSegment = self.index.getSegment(self.readState.currentSeq);
  101. if (self.readState.currentSegment) {
  102. var url = self.readState.currentSegment.uri;
  103. function tryfetch(start) {
  104. var seq = self.readState.currentSeq;
  105. fetchfrom(seq, self.readState.currentSegment, start, function(err) {
  106. if (err) {
  107. if (!self.keepConnection) return self.emit('error', err);
  108. console.error('While fetching '+url+':', err.stack || err);
  109. // retry with missing range if it is still relevant
  110. if (err instanceof uristream.PartialError && err.processed > 0 &&
  111. self.index.getSegment(seq))
  112. return tryfetch(start + err.processed);
  113. }
  114. self.readState.currentSegment = null;
  115. if (seq === self.readState.currentSeq)
  116. self.readState.currentSeq++;
  117. checkcurrent();
  118. });
  119. }
  120. tryfetch(0);
  121. } else if (self.index.ended)
  122. self.end();
  123. else if (!self.index.type && (self.index.lastSeqNo() < self.readState.currentSeq-1)) {
  124. // handle live stream restart
  125. self.readState.currentSeq = self.index.startSeqNo(true);
  126. checkcurrent();
  127. }
  128. }
  129. function fetchfrom(seqNo, segment, start, cb) {
  130. var segmentUrl = url.resolve(self.baseUrl, segment.uri)
  131. var probe = !!self.noData;
  132. debug('fetching segment', segmentUrl);
  133. var stream = uristream(segmentUrl, { probe:probe, start:start, highWaterMark:100*1000*1000 });
  134. stream.on('meta', function(meta) {
  135. debug('got segment info', meta);
  136. if (meta.mime !== 'video/mp2t'/* &&
  137. meta.mime !== 'audio/aac' && meta.mime !== 'audio/x-aac' &&
  138. meta.mime !== 'audio/ac3'*/) {
  139. if (stream.abort) stream.abort();
  140. return stream.emit(new Error('Unsupported segment MIME type: '+meta.mime));
  141. }
  142. // abort() indicates a temporal stream. Ie. ensure it is completed in a timely fashion
  143. if (self.index.isLive() && typeof stream.abort == 'function') {
  144. var duration = segment.duration || self.index.target_duration || 10;
  145. duration = Math.min(duration, self.index.target_duration || 10);
  146. self.readState.timer = setTimeout(function() {
  147. if (self.readState.stream) {
  148. debug('timed out waiting for data');
  149. self.readState.stream.abort();
  150. }
  151. // TODO: ensure Done() is always called
  152. self.readState.timer = null;
  153. }, 1.5*duration*1000);
  154. }
  155. if (start === 0)
  156. self.emit('segment', seqNo, segment.duration, meta);
  157. });
  158. stream.pipe(self);
  159. oncemore(stream).once('end', 'error', function(err) {
  160. clearTimeout(self.readState.timer);
  161. stream.unpipe(self);
  162. self.readState.stream = null;
  163. if (err) debug('stream error', err);
  164. else debug('finished with input stream', stream.meta.url);
  165. cb(err);
  166. });
  167. self.readState.stream = stream;
  168. }
  169. // allow piping content to self
  170. this.write = this.push.bind(this);
  171. this.end = function() {};
  172. Readable.call(this, options);
  173. }
  174. util.inherits(HlsStreamReader, Readable);
  175. HlsStreamReader.prototype._read = function(n, cb) {
  176. this.emit('drain');
  177. };
  178. function hlsreader(url, options) {
  179. return new HlsStreamReader(url, options);
  180. }