From 4f5fb590bf5d429c76229cddffae47fe77fbc061 Mon Sep 17 00:00:00 2001 From: Alexandre Storelli Date: Thu, 18 Jan 2018 00:12:09 +0100 Subject: [PATCH] server works with delay for audio, delay progressbar in UI --- DlFactory.js | 170 +++++++++++++++++++++++-- client/package.json | 3 +- client/src/App.js | 92 ++++++++------ client/src/DelayCanvas.js | 88 +++++++++++++ client/src/Metadata.js | 15 +-- client/src/load.js | 5 +- config.js | 15 +++ findDataFiles.js | 139 +++++++++++++++++++++ handlers.js | 234 ++++++++++++++++++++++++++++++++++ index.js | 255 ++++++-------------------------------- 10 files changed, 742 insertions(+), 274 deletions(-) create mode 100644 client/src/DelayCanvas.js create mode 100644 config.js create mode 100644 findDataFiles.js create mode 100644 handlers.js diff --git a/DlFactory.js b/DlFactory.js index 967a4d9..20f7b38 100644 --- a/DlFactory.js +++ b/DlFactory.js @@ -7,7 +7,7 @@ var cp = require("child_process"); var fs = require("fs"); var { getMeta } = require("webradio-metadata"); var Dl = require("../adblockradio-dl/dl.js"); - +var config = require("./config.js"); class Db { constructor(options) { @@ -15,11 +15,17 @@ class Db { this.name = options.name; this.path = options.path; this.ext = options.ext; + this.audioCache = new AudioCache({ bitrate: options.bitrate, cacheLen: config.user.cacheLen }); + this.metaCache = new MetaCache({ cacheLen: config.user.cacheLen }); + } + + dirDate(now) { + return (now.getUTCFullYear()) + "-" + (now.getUTCMonth()+1 < 10 ? "0" : "") + (now.getUTCMonth()+1) + "-" + (now.getUTCDate() < 10 ? "0" : "") + (now.getUTCDate()); } newAudioSegment() { var now = new Date(); - var dir = this.path + "/records/" + dirDate(now) + "/" + this.country + "_" + this.name + "/todo/"; + var dir = this.path + "/records/" + this.dirDate(now) + "/" + this.country + "_" + this.name + "/todo/"; var path = dir + now.toISOString(); log.debug("newAudioSegment: path=" + path); var self = this; @@ -32,6 +38,7 @@ class Db { return { audio: new AudioWriteStream(path + "." + self.ext), //new fs.createWriteStream(path + "." + self.ext), // metadata: new MetaWriteStream(path + ".json"), + date: now, path: path }; } @@ -71,6 +78,67 @@ class AudioWriteStream extends Duplex { } } +class AudioCache extends Writable { + constructor(options) { + super(); + this.buffer = null; + this.cacheLen = options.cacheLen; + this.bitrate = options.bitrate; + this.flushAmount = this.cacheLen * this.bitrate * 0.1; + this.readCursor = null; + } + + _write(data, enc, next) { + this.buffer = this.buffer ? Buffer.concat([ this.buffer, data ]) : data; + //log.debug("AudioCache: _write: add " + data.length + " to buffer, new len=" + this.buffer.length); + if (this.buffer.length >= this.cacheLen * this.bitrate + this.flushAmount) { + log.debug("AudioCache: _write: cutting buffer at len = " + this.cacheLen * this.bitrate); + this.buffer = this.buffer.slice(this.cacheLen * this.bitrate); + if (this.readCursor) { + this.readCursor -= this.flushAmount; + if (this.readCursor <= 0) this.readCursor = null; + } + } + next(); + } + + readLast(secondsFromEnd, duration) { + var l = this.buffer.length; + if (secondsFromEnd < 0 || duration < 0) { + log.error("AudioCache: readLast: negative secondsFromEnd or duration"); + return null; + } else if (duration > secondsFromEnd) { + log.error("AudioCache: readLast: duration=" + duration + " higher than secondsFromEnd=" + secondsFromEnd); + return null; + } else if (secondsFromEnd * this.bitrate >= l) { + log.error("AudioCache: readLast: attempted to read " + secondsFromEnd + " seconds (" + secondsFromEnd * this.bitrate + " b) while bufferLen=" + l); + return null; + } + var data; + if (duration) { + data = this.buffer.slice(l - secondsFromEnd * this.bitrate, l - (secondsFromEnd-duration) * this.bitrate); + this.readCursor = l - (secondsFromEnd-duration) * this.bitrate; + } else { + data = this.buffer.slice(l - secondsFromEnd * this.bitrate); + this.readCursor = l; + } + return data; + } + + readAmountAfterCursor(duration) { + var nextCursor = this.readCursor + duration * this.bitrate; + if (duration < 0) { + log.error("AudioCache: readAmountAfterCursor: negative duration"); + return null; + } else if (nextCursor >= this.buffer.length) { + log.warn("AudioCache: readAmountAfterCursor: will read until " + this.buffer.length + " instead of " + nextCursor); + } + var data = this.buffer.slice(this.readCursor, nextCursor); + this.readCursor = Math.min(this.buffer.length, nextCursor); + return data; + } +} + class MetaWriteStream extends Writable { constructor(path) { super({ objectMode: true }); @@ -97,11 +165,90 @@ class MetaWriteStream extends Writable { } } -var dirDate = function(now) { - return (now.getUTCFullYear()) + "-" + (now.getUTCMonth()+1 < 10 ? "0" : "") + (now.getUTCMonth()+1) + "-" + (now.getUTCDate() < 10 ? "0" : "") + (now.getUTCDate()); +class MetaCache extends Writable { + constructor(options) { + super({ objectMode: true }); + this.meta = {}; + this.cacheLen = options.cacheLen; + } + + _write(meta, enc, next) { + if (!meta.type) { + log.error("MetaCache: no data type"); + return next(); + } else { + //log.debug("MetaCache: _write: " + JSON.stringify(meta)); + } + // events of this kind: + // meta = { type: "metadata", validFrom: Date, validTo: Date, payload: { artist: "...", title : "...", cover: "..." } } + // meta = { type: "class", validFrom: Date, validTo: Date, payload: "todo" } + // meta = { type: "signal", validFrom: Date, validTo: Date, payload: [0.4, 0.3, ...] } + + // are stored in the following structure: + // this.meta = { + // "metadata": [ + // { validFrom: ..., validTo: ..., payload: { ... } }, (merges the contiguous segments) + // ... + // ], + // "class": [ + // { validFrom: ..., validTo: ..., payload: ... }, (merges the contiguous segments) + // ... + // ], + // "signal": [ + // { validFrom: ..., validTo: ..., payload: [ ... ] }, + // ... + // ] + // } + + switch (meta.type) { + case "metadata": + case "class": + if (!this.meta[meta.type]) { + this.meta[meta.type] = [ { validFrom: meta.validFrom, validTo: meta.validTo, payload: meta.payload } ]; + } else { + var samePayload = true; + for (var key in meta.payload) { + if (meta.payload[key] !== this.meta[meta.type][this.meta[meta.type].length-1].payload[key]) { + samePayload = false; + //log.debug("MetaCache: _write: different payload key=" + key + " new=" + meta.payload[key] + " vs old=" + this.meta[meta.type][this.meta[meta.type].length-1].payload[key]); + break; + } + } + if (samePayload) { + this.meta[meta.type][this.meta[meta.type].length-1].validTo = meta.validTo; // extend current segment validity + } else { + this.meta[meta.type][this.meta[meta.type].length-1].validTo = meta.validFrom; // create a new segment + this.meta[meta.type].push({ validFrom: meta.validFrom, validTo: meta.validTo, payload: meta.payload }); + } + } + break; + case "signal": + if (!this.meta[meta.type]) { + this.meta[meta.type] = [ { validFrom: meta.validFrom, validTo: meta.validTo, payload: meta.payload } ]; + } else { + this.meta[meta.type].push({ validFrom: meta.validFrom, validTo: meta.validTo, payload: meta.payload }); + } + break; + default: + log.error("MetaCache: _write: unknown metadata type = " + meta.type); + } + + // clean old entries + while (+this.meta[meta.type][0].validTo <= +new Date() - 1000 * this.cacheLen) { + this.meta[meta.type].splice(0, 1); + } + + //log.debug("MetaCache: _write: meta[" + meta.type + "]=" + JSON.stringify(this.meta[meta.type])); + next(); + } + + read() { + this.meta.now = +new Date(); + return this.meta; + } } -var DlFactory = function(radio, options) { +module.exports = function(radio, options) { var newDl = new Dl({ country: radio.country, name: radio.name, segDuration: options.segDuration }); var dbs = null; newDl.on("error", function(err) { @@ -111,7 +258,7 @@ var DlFactory = function(radio, options) { //metadataCallback(metadata); radio.url = metadata.url; radio.favicon = metadata.favicon; - var db = new Db({ country: radio.country, name: radio.name, ext: metadata.ext, path: __dirname }); + var db = new Db({ country: radio.country, name: radio.name, ext: metadata.ext, bitrate: metadata.bitrate, path: __dirname }); log.debug("DlFactory: " + radio.country + "_" + radio.name + " metadata=" + JSON.stringify(metadata)); newDl.on("data", function(dataObj) { @@ -128,7 +275,9 @@ var DlFactory = function(radio, options) { dbs = db.newAudioSegment(); Object.assign(radio.liveStatus, { currentPrefix: dbs.path, - liveReadStream: dbs.audio + liveReadStream: dbs.audio, + audioCache: db.audioCache, + metaCache: db.metaCache }); if (options.fetchMetadata) { getMeta(radio.country, radio.name, function(err, parsedMeta, corsEnabled) { @@ -138,7 +287,9 @@ var DlFactory = function(radio, options) { Object.assign(radio.liveStatus, { metadata: parsedMeta }); - dbs.metadata.write({ type: "title", data: parsedMeta }); + dbs.metadata.write({ type: "metadata", data: parsedMeta }); + db.metaCache.write({ type: "metadata", validFrom: +dbs.date-500*options.segDuration, validTo: +dbs.date+500*options.segDuration, payload: parsedMeta }); + db.metaCache.write({ type: "class", validFrom: +dbs.date-500*options.segDuration, validTo: +dbs.date+500*options.segDuration, payload: "todo" }); } else { log.warn("getMeta: could not write metadata, stream already ended"); } @@ -147,9 +298,8 @@ var DlFactory = function(radio, options) { dbs.audio.write(dataObj.data); newDl.resume(); } + db.audioCache.write(dataObj.data); }); }); return newDl; } - -module.exports = DlFactory; diff --git a/client/package.json b/client/package.json index a92de0e..a88af4a 100644 --- a/client/package.json +++ b/client/package.json @@ -17,5 +17,6 @@ "build": "react-scripts build", "test": "react-scripts test --env=jsdom", "eject": "react-scripts eject" - } + }, + "homepage" : "https://www.adblockradio.com/buffer" } diff --git a/client/src/App.js b/client/src/App.js index eaadc95..92fd7dc 100644 --- a/client/src/App.js +++ b/client/src/App.js @@ -4,6 +4,7 @@ import React, { Component } from 'react'; import './App.css'; import Radio from './Radio.js'; import Metadata from './Metadata.js'; +import DelayCanvas from './DelayCanvas.js'; import { load, refreshMetadata, HOST } from './load.js'; import { play, stop } from './audio.js'; @@ -16,8 +17,6 @@ import iconPlay from "./img/start_1279169.svg"; import iconStop from "./img/stop_1279170.svg"; //import iconNext from "./img/next_607554.svg"; -const LIVE_DELAY = 15000; - class App extends Component { constructor(props) { super(props); @@ -25,7 +24,8 @@ class App extends Component { configLoaded: false, config: [], playingRadio: null, - playingDate: null + playingDate: null, + clockDiff: 0 } this.play = this.play.bind(this); this.seekBackward = this.seekBackward.bind(this); @@ -40,10 +40,15 @@ class App extends Component { var f = function(iradio, callback) { if (iradio >= self.state.config.radios.length) return callback(); var radio = self.state.config.radios[iradio].country + "_" + self.state.config.radios[iradio].name; - var startDate = self.state.playingRadio === radio ? new Date(+new Date() - 120*60000) : new Date(+new Date() - 60000); - refreshMetadata(radio, startDate.toISOString(), function(metadata) { - metadata[metadata.length-1].end = null; - stateChange["meta" + radio] = metadata.reverse(); + refreshMetadata(radio, function(metadata) { + for (var type in metadata) { + if (type === "now") { + stateChange.clockDiff = +new Date() - metadata.now; + } else { + metadata[type][metadata[type].length-1].validTo = null; + stateChange[radio + "|" + type] = metadata[type].reverse(); + } + } f(iradio+1, callback); }); } @@ -52,7 +57,7 @@ class App extends Component { }); } - load("/config", function(res) { + load("/config?t=" + Math.round(Math.random()*1000000), function(res) { try { var config = JSON.parse(res); console.log(config); @@ -77,32 +82,30 @@ class App extends Component { this.setState({ date: new Date() }); } - play(radio, date) { - var self = this; - if (radio) { - console.log("playing delta = " + (+new Date() - new Date(date))); - if (!date) { - date = new Date(+new Date() - LIVE_DELAY).toISOString(); - this.setState({ playingLive: true }); - } else { - this.setState({ playingLive: false }); - } - + play(radio, delay) { + if (radio || delay) { + //var delay = +new Date() - this.state.clockDiff - (date ? new Date(date) : new Date(); + //var secondsDelay = Math.round(delay/1000); + if (!radio) radio = this.state.playingRadio; + if (!delay) delay = 0; + console.log("Play: radio=" + radio + " delay=" + delay); this.setState({ playingRadio: radio, - playingDate: date, - playingDelta: +new Date() - new Date(date) + //playingDate: date, + playingDelay: delay, + playingLive: delay === 0 }); - play(HOST + "/listen/" + radio + "/" + date); + + play(HOST + "/listen/" + radio + "/" + Math.round(delay/1000)); document.title = radio.split("_")[1] + " - Adblock Radio"; } else { this.setState({ playingRadio: null, - playingDate: null, - playingDelta: 0, + //playingDate: null, + playingDelay: null, playingLive: null }); - clearInterval(self.metadataTimerID); + //clearInterval(self.metadataTimerID); stop(); document.title = "Adblock Radio"; } @@ -110,17 +113,18 @@ class App extends Component { seekBackward() { if (!this.state.playingRadio) return; - this.play(this.state.playingRadio, new Date(+new Date(this.state.playingDate) - 30000).toISOString()); + this.play(this.state.playingRadio, Math.min(this.state.playingDelay + 30000, this.state.config.user.cacheLen*1000)); } seekForward() { if (!this.state.playingRadio) return; - var targetDate = +new Date(this.state.playingDate) + 30000; - if (targetDate > +new Date() - LIVE_DELAY) { // switch to live + //var targetDate = +new Date(this.state.playingDate) + 30000; + /*if (targetDate > +new Date() - LIVE_DELAY) { // switch to live this.play(this.state.playingRadio); } else { this.play(this.state.playingRadio, new Date(targetDate).toISOString()); - } + }*/ + this.play(this.state.playingRadio, Math.max(this.state.playingDelay - 30000,0)); } render() { @@ -144,14 +148,17 @@ class App extends Component { var statusText; if (self.state.playingRadio) { + var delayText = " (direct)"; + if (!self.state.playingLive) { + var delaySeconds = Math.round(self.state.playingDelay/1000); // + self.state.config.user.streamInitialBuffer); + var delayMinutes = Math.floor(delaySeconds / 60); + delaySeconds = delaySeconds % 60; + delayText = " (en différé de " + (delayMinutes ? delayMinutes + "m" : "") + (delaySeconds < 10 ? "0" : "") + delaySeconds + "s)"; + } statusText = ( {self.state.playingRadio.split("_")[1]} - {self.state.playingLive ? - " (direct)" - : - " (" + moment(+self.state.date - self.state.playingDelta + LIVE_DELAY).fromNow() + ")" - } + {delayText} ) } else { @@ -177,6 +184,7 @@ class App extends Component { ); + //console.log("Metadata props: date=" + (+self.state.date) + " clockDiff=" + self.state.clockDiff + " playingDelay=" + self.state.playingDelay); return ( @@ -189,22 +197,26 @@ class App extends Component { playCallback={self.play} playing={playing} showMetadata={!self.state.playingRadio} - liveMetadata={self.state["meta" + radio.country + "_" + radio.name]} + liveMetadata={self.state[radio.country + "_" + radio.name + "|metadata"]} key={"radio" + i} /> ) })} {self.state.playingRadio && } {status} {buttons} + {/**/} @@ -249,7 +261,7 @@ const Controls = styled.div` `; const StatusTextContainer = styled.span` - padding: 0 20px; + padding: 10px 20px 0 20px; `; const StatusClock = styled.span` @@ -257,7 +269,7 @@ const StatusClock = styled.span` `; const StatusButtonsContainer = styled.span` - + padding: 10px 0 0 0; `; const PlaybackButton = styled.img` diff --git a/client/src/DelayCanvas.js b/client/src/DelayCanvas.js new file mode 100644 index 0000000..743b130 --- /dev/null +++ b/client/src/DelayCanvas.js @@ -0,0 +1,88 @@ +// Copyright (c) 2018 Alexandre Storelli + +import React, { Component } from "react"; +import PropTypes from "prop-types"; +//import styled from "styled-components"; +import classNames from 'classnames'; + + +class DelayCanvas extends Component { + + constructor(props) { + super(props); + this.updateCanvas = this.updateCanvas.bind(this); + this.play = this.play.bind(this); + this.getCursorPosition = this.getCursorPosition.bind(this); + + } + + componentDidMount() { + this.updateHandle = setInterval(this.updateCanvas, 2000); + this.refs.canvas.addEventListener("mousedown", this.getCursorPosition); + } + + componentWillUnmount() { + clearInterval(this.updateHandle); + //window.cancelAnimationFrame(this.updateHandle); + } + + componentDidUpdate() { + //this.drawBg(); + this.updateCanvas(); + } + + play(delay) { + this.props.playCallback(null, delay); + } + + getCursorPosition(event) { + var rect = this.refs.canvas.getBoundingClientRect(); + var x = event.clientX - rect.left; + + var canvasDom = document.getElementById('canvas'); + var cs = getComputedStyle(canvasDom); + var width = parseInt(cs.getPropertyValue('width'), 10); + + //var width = this.refs.canvas.getContext("2d").canvas.width; + var newDelay = Math.round(this.props.cacheLen*(1-x/width)*1000); + //console.log("Canvas click: x=" + x + " width=" + width + " cacheLen=" + this.props.cacheLen + " newDelay=" + newDelay); + this.play(newDelay); + } + + updateCanvas() { + const ctx = this.refs.canvas.getContext("2d"); + let height = ctx.canvas.height; //parseFloat(this.props.style.height); + let width = ctx.canvas.width; //parseFloat(this.props.style.width); + ctx.clearRect(0, 0, width, height); + ctx.fillStyle = "rgba(0,0,128,1)"; + //console.log("canvas height=" + height + " px width=" + width + "px playingDelay=" + this.props.playingDelay + " cacheLen=" + this.props.cacheLen); + ctx.fillRect(0, 0, Math.round(width*(1-this.props.playingDelay/1000/this.props.cacheLen)), height); + } + + render() { + return ( + + ) + } +} + +DelayCanvas.propTypes = { + playing: PropTypes.bool.isRequired, + playingDelay: PropTypes.number, + cacheLen: PropTypes.number.isRequired, + playCallback: PropTypes.func.isRequired +}; + +/*var canvasContainerStyle = { + width: "100%", + height: "10px" +}*/ + +var canvasStyle = { + position: "absolute", + width: "100%", + height: "10px", + alignSelf: "flex-start", +} + +export default DelayCanvas; diff --git a/client/src/Metadata.js b/client/src/Metadata.js index 204cd86..a20ebc9 100644 --- a/client/src/Metadata.js +++ b/client/src/Metadata.js @@ -19,7 +19,8 @@ class Metadata extends Component { play(when) { - this.props.playCallback(this.props.playingRadio, when); + console.log("Metadata play: date=" + this.props.date + " when=" + when + " delay=" + (this.props.date - when)); + this.props.playCallback(this.props.playingRadio, this.props.date - when); } stop() { @@ -33,16 +34,16 @@ class Metadata extends Component { {this.props.metaList ? this.props.metaList.map(function(item, i) { - if (!item.title || (self.props.maxItems && i >= self.props.maxItems)) return null + if (!item.payload || (self.props.maxItems && i >= self.props.maxItems)) return null //var playerTime = +self.props.playingDate; - var playing = +new Date(item.start) - 1000 <= self.props.playingDate && (!item.end || self.props.playingDate < +new Date(item.end) - 1000); + var playing = +item.validFrom - 1000 <= (self.props.date-self.props.playingDelay) && (!item.validTo || (self.props.date-self.props.playingDelay) < +item.validTo - 1000); //console.log("playing=" + playing + " start=" + item.start + " date=" + playerTime + " stop=" + item.end); return ( - + - {!previewMode && (moment(item.start).format("HH:mm") + " – ")}{item.title.artist} - {item.title.title} + {!previewMode && (moment(item.validFrom).format("HH:mm") + " – ")}{item.payload.artist} - {item.payload.title} - + ) }) @@ -79,7 +80,7 @@ const MetadataItem = styled.div` background: #eee; display: flex; cursor: pointer; - + &.playing { border: 2px solid red; } diff --git a/client/src/load.js b/client/src/load.js index a504648..2083cb9 100644 --- a/client/src/load.js +++ b/client/src/load.js @@ -1,4 +1,5 @@ const HOST = "http://localhost:9820"; +//const HOST = "https://bufferapi.s00.adblockradio.com" exports.load = function(path, callback) { var xhttp = new XMLHttpRequest(); @@ -16,8 +17,8 @@ exports.load = function(path, callback) { exports.HOST = HOST; -exports.refreshMetadata = function(radio, date, callback) { - exports.load("/metadata/" + radio + "/" + date, function(res) { +exports.refreshMetadata = function(radio, callback) { + exports.load("/metadata/" + radio + "/0", function(res) { var metadata = []; try { metadata = JSON.parse(res); diff --git a/config.js b/config.js new file mode 100644 index 0000000..babb132 --- /dev/null +++ b/config.js @@ -0,0 +1,15 @@ +var log = require("loglevel"); +var fs = require("fs"); + +// list of listened radios: +var config = new Object(); +try { + var radiosText = fs.readFileSync("config/radios.json"); + config.radios = JSON.parse(radiosText); + var userText = fs.readFileSync("config/user.json"); + config.user = JSON.parse(userText); +} catch(e) { + return log.error("cannot load config. err=" + e); +} + +module.exports = config; diff --git a/findDataFiles.js b/findDataFiles.js new file mode 100644 index 0000000..bb5c206 --- /dev/null +++ b/findDataFiles.js @@ -0,0 +1,139 @@ +var fs = require("fs"); +var log = require("loglevel"); +log.setLevel("debug"); +var async = require("async"); +var consts = { + WLARRAY: ["0-ads", "1-speech", "2-music", "9-unsure", "mrs", "todo"] +} + +var getDirs = function(rootDir, cb) { + fs.readdir(rootDir, function(err, files) { + var dirs = []; + if (err) { + log.warn("getDirs: readdir error for " + rootDir + " err=" + err); + return cb(dirs); + } + //log.debug("getDirs: files found in " + rootDir + " : " + files); + + var f = function(i, callback) { + if (i >= files.length) return callback(); + + var filePath = rootDir + '/' + files[i]; + fs.stat(filePath, function(err, stat) { + if (err) { + log.warn("getDirs: stat error for " + filePath + " err=" + err); + } + if (stat.isDirectory()) { + dirs.push(files[i]); + } + f(i+1, callback); + }); + } + + f(0, function() { + return cb(dirs); + }); + }); +} + +var getFiles = function(path, after, before, cb) { + fs.readdir(path, function (err, files) { + if (err) { + log.warn("error listing files in path " + path + ". err=" + err); + return cb([]); + } + var list = {}; + + // only get files that have good date + for (var i=files.length-1; i>=0; i--) { + if ((after && files[i] < after) || // file names begin by ISO dates (24 characters) + (before && files[i] > before)) { + //log.debug("getFiles: remove " + files[i]); + files.splice(i, 1); + } else { + var spf = files[i].split("Z."); // end of the ISO date + if (spf.length != 2) log.warn("getFiles: malformed file: " + files[i]); + var path1 = path + "/" + spf[0] + "Z"; + if (!list[path1]) { + var ps = path.split("/"); + list[path1] = { class: ps[ps.length-1] }; + } + /*if (spf[1].slice(-5) == ".part") { + list[path1]["partial"] = true; + list[path1][spf[1].slice(spf[1].length-5)] = true; + } else {*/ + list[path1][spf[1]] = true; + //} + } + } + cb(list); + }); +} + +// finds all files in the subdirectory structure +// options: { +// radios: [country1_name1, country2_name2, ...] OR country: ... name: ... +// before/after: ISO Date, e.g. new Date().toISOString() +// path: __dirname/records/DATE/RADIO/CLASS/ISODATE.* +// } +var findDataFiles = function(options, callback) { + var targetRadios = (options && options.radios) ? options.radios : [options.country + "_" + options.name]; + var timeFrame = { + after: (options && options.after) ? options.after : null, + before: (options && options.before) ? options.before : null + }; + + var files = new Object(); + for (let i=0; i=0; i--) { + if (timeFrame.after && dateDirs[i] < timeFrame.after.slice(0,10)) dateDirs.splice(i, 1); + if (timeFrame.before && dateDirs[i] > timeFrame.before.slice(0,10)) dateDirs.splice(i, 1); + } + log.debug("findDataFiles: dateDirs:" + JSON.stringify(dateDirs)); + + async.forEachOf(dateDirs, function(dateDir, index, dateDirCallback) { + getDirs(options.path + "/records/" + dateDir, function(radioDirs) { + //log.debug("findDataFiles: radioDirs before = " + radioDirs); + for (let i=radioDirs.length-1; i>=0; i--) { + if (targetRadios.indexOf(radioDirs[i]) < 0) radioDirs.splice(i, 1); + } + log.debug("findDataFiles: radioDirs=" + radioDirs); + + async.forEachOf(radioDirs, function(radioDir, index, radioDirCallback) { + async.forEachOf(consts.WLARRAY, function(dataDir, index, dataDirCallback) { + var path = options.path + "/records/" + dateDir + "/" + radioDir + "/" + dataDir; + fs.stat(path, function(err, stat) { + if (stat && stat.isDirectory()) { + getFiles(path, timeFrame.after, timeFrame.before, function(partialFiles) { + log.debug("findDataFiles: " + dataDir + ": " + Object.keys(partialFiles).length + " files found"); + Object.assign(files[dataDir], partialFiles); // = files[dataDir].concat(fullPathFiles); + dataDirCallback(); + }); + } else { + dataDirCallback(); + } + }); + + }, function(err) { + if (err) log.error("findDataFiles: pb during data dir listing. err=" + err.message); + radioDirCallback(); + }); + + }, function(err) { + if (err) log.error("findDataFiles: pb during radio dir listing. err=" + err.message); + dateDirCallback(); + }); + }); + + }, function(err) { + if (err) log.error("findDataFiles: pb during date dir listing. err=" + err.message); + callback(files); + }); + }); +} + +module.exports = findDataFiles; diff --git a/handlers.js b/handlers.js new file mode 100644 index 0000000..a2945c6 --- /dev/null +++ b/handlers.js @@ -0,0 +1,234 @@ +"use strict"; + +var config = require("./config.js"); +var log = require("loglevel"); + +var getDateFromPath = function(path) { + var spl = path.split("/"); + return spl[spl.length-1]; +} + +var getRadio = function(country, name) { + if (name) { // both parameters used + for (var j=0; j 0) { + state.nSegmentsInitialBuffer--; + filesCallback(); + } else { + var delay = +SEG_DURATION*1000 - (new Date()-state.lastSentDate); + log.debug("listenHandler: send more data in " + delay + " ms"); + setTimeout(filesCallback, delay); + } + } + }, function(err) { + if (err) { + log.error("listen: err=" + err.message); + callback(); + } else { + listenHandler(response, radio, before, state, callback); + } + }); + }); +}*/ + +//const STREAM_INITIAL_BUFFER = 8; // send N seconds at connection. +//const STREAM_GRANULARITY = 1; // send N seconds of data every N seconds. +var listenRequestDate = null; +exports.listenHandler = function(response, radio, delay, state, callback) { + var initialBuffer = getRadio(radio).liveStatus.audioCache.readLast(+delay+config.user.streamInitialBuffer,config.user.streamInitialBuffer); + log.info("listen: send initial buffer of " + initialBuffer.length + " bytes"); + response.write(initialBuffer); + + var listenTimer = setInterval(function() { + if (state.newRequest) { + listenRequestDate = state.requestDate; + state.newRequest = false; + } else if (listenRequestDate !== state.requestDate) { + log.warn("request canceled because another one has been initiated"); + clearInterval(listenTimer); + return callback(); + } + var willWaitDrain = !response.write(""); + if (!willWaitDrain) { // detect congestion of stream + sendMore(); + } else { + log.debug("listenHandler: will wait for drain event"); + + var drainCallback = function() { + clearTimeout(timeoutMonitor); + sendMore(); + } + response.once("drain", drainCallback); + var timeoutMonitor = setTimeout(function() { + response.removeListener("drain", drainCallback); + clearInterval(listenTimer); + log.error("listenHandler: drain event not emitted, connection timeout"); + callback(); + }, config.user.streamGranularity*1500); + } + }, 1000*config.user.streamGranularity); + var sendMore = function() { + response.write(getRadio(radio).liveStatus.audioCache.readAmountAfterCursor(config.user.streamGranularity)); + } +} + +/*exports.metadataHandler = function(response, radio, after) { + findDataFiles({ radios: [ radio ], after: after, path: __dirname }, function(classes) { + //log.debug(classes); // very verbose + var list = {}; + for (var classItem in classes) { + Object.assign(list, classes[classItem]); + } + var files = Object.keys(list); + files.sort(); + //response.json(list); + var isSameSegment = function(curSegment, newData) { + if (curSegment.class !== newData.class) { + //log.debug("isSameSegment: curEnd=" + curSegment.end + " diff class cur=" + curSegment.class + " vs new=" + newData.class); + return false; + } else if (!curSegment.title || !newData.title) { + //log.debug("isSameSegment: curEnd=" + curSegment.end + " no title"); + return true; + } else if (curSegment.title.artist && newData.title.artist && curSegment.title.artist !== newData.title.artist) { + //log.debug("isSameSegment: curEnd=" + curSegment.end + " diff artist cur=" + curSegment.title.artist + " vs new=" + newData.title.artist); + return false; + } else if (curSegment.title.title && newData.title.title && curSegment.title.title !== newData.title.title) { + //log.debug("isSameSegment: curEnd=" + curSegment.end + " diff title cur=" + curSegment.title.title + " vs new=" + newData.title.title); + return false; + } else if (curSegment.title.cover && newData.title.cover && curSegment.title.cover !== newData.title.cover) { + //log.debug("isSameSegment: curEnd=" + curSegment.end + " diff cover"); + return false; + } else { + //log.debug("isSameSegment: same segment"); + return true; + } + } + var newSegment = function(index) { + return { class: list[files[index]].class, start: getDateFromPath(files[index]), end: getDateFromPath(files[index]), title: null } + } + + var result = []; + if (files.length == 0) return response.json(result); + var curSegment = newSegment(0); + async.forEachOfSeries(files, function(file, index, filesCallback) { + //log.debug("metadata: read=" + file + ".json"); + + curSegment.end = getDateFromPath(file); + + if (isLive(radio, file)) { + //log.info("listenHandler: currentAudioFile"); + var pData = { class: list[file].class, title: getRadio(radio).liveStatus.metadata }; + if (!isSameSegment(curSegment, pData)) { + result.push(curSegment); + curSegment = newSegment(index); + } + curSegment.title = pData.title; + filesCallback(); + } else { + fs.readFile(file + ".json", function(err, data) { + if (err) log.error("metadata: readFile err=" + err); + try { + var pData = JSON.parse(data); + } catch(e) { + log.error("metadata: json parse error file=" + file + ".json"); + return filesCallback(); + } + pData.class = list[file].class; + if (!isSameSegment(curSegment, pData)) { + result.push(curSegment); + curSegment = newSegment(index); + } + curSegment.title = pData.title; + filesCallback(); + }); + } + + + }, function(err) { + if (err) log.error("metadata: err=" + err.message); + result.push(curSegment); + response.set({ 'Access-Control-Allow-Origin': '*' }); + response.json(result); + }); + }); +}*/ diff --git a/index.js b/index.js index 89cf683..6b683f3 100644 --- a/index.js +++ b/index.js @@ -4,29 +4,19 @@ "use strict"; var abrsdk = require("adblockradio-sdk"); -var fs = require("fs"); var log = require("loglevel"); log.setLevel("debug"); var cp = require("child_process"); -var findDataFiles = require("../adblockradio/predictor-db/findDataFiles.js"); +var findDataFiles = require("./findDataFiles.js"); var async = require("async"); var DlFactory = require("./DlFactory.js"); const DL = true; const FETCH_METADATA = true; -const SEG_DURATION = 10; -const LISTEN_BUFFER = 30; +const SEG_DURATION = 10; // in seconds +const LISTEN_BUFFER = 30; // in seconds -// list of listened radios: -var config = new Object(); -try { - var radiosText = fs.readFileSync("config/radios.json"); - config.radios = JSON.parse(radiosText); - var userText = fs.readFileSync("config/user.json"); - config.user = JSON.parse(userText); -} catch(e) { - return log.error("cannot load config. err=" + e); -} +var config = require("./config.js"); if (DL) { var dl = []; @@ -46,234 +36,71 @@ server.listen(9820, "localhost"); app.get('/config', function(request, response) { response.set({ 'Access-Control-Allow-Origin': '*' }); - response.json(config); + var result = { radios: [], user: config.user }; + for (var i=0; i 0) { - state.nSegmentsInitialBuffer--; - filesCallback(); - } else { - var delay = +SEG_DURATION*1000 - (new Date()-state.lastSentDate); - log.debug("listenHandler: send more data in " + delay + " ms"); - setTimeout(filesCallback, delay); - } - } - }, function(err) { - if (err) { - log.error("listen: err=" + err.message); - callback(); - } else { - listenHandler(response, radio, before, state, callback); - } - }); - }); -} - -var metadataHandler = function(response, radio, after) { - findDataFiles({ radios: [ radio ], after: after, path: __dirname }, function(classes) { - //log.debug(classes); // very verbose - var list = {}; - for (var classItem in classes) { - Object.assign(list, classes[classItem]); - } - var files = Object.keys(list); - files.sort(); - //response.json(list); - var isSameSegment = function(curSegment, newData) { - if (curSegment.class !== newData.class) { - //log.debug("isSameSegment: curEnd=" + curSegment.end + " diff class cur=" + curSegment.class + " vs new=" + newData.class); - return false; - } else if (!curSegment.title || !newData.title) { - //log.debug("isSameSegment: curEnd=" + curSegment.end + " no title"); - return true; - } else if (curSegment.title.artist && newData.title.artist && curSegment.title.artist !== newData.title.artist) { - //log.debug("isSameSegment: curEnd=" + curSegment.end + " diff artist cur=" + curSegment.title.artist + " vs new=" + newData.title.artist); - return false; - } else if (curSegment.title.title && newData.title.title && curSegment.title.title !== newData.title.title) { - //log.debug("isSameSegment: curEnd=" + curSegment.end + " diff title cur=" + curSegment.title.title + " vs new=" + newData.title.title); - return false; - } else if (curSegment.title.cover && newData.title.cover && curSegment.title.cover !== newData.title.cover) { - //log.debug("isSameSegment: curEnd=" + curSegment.end + " diff cover"); - return false; - } else { - //log.debug("isSameSegment: same segment"); - return true; - } - } - var newSegment = function(index) { - return { class: list[files[index]].class, start: getDateFromPath(files[index]), end: getDateFromPath(files[index]), title: null } - } - - var result = []; - if (files.length == 0) return response.json(result); - var curSegment = newSegment(0); - async.forEachOfSeries(files, function(file, index, filesCallback) { - //log.debug("metadata: read=" + file + ".json"); - - curSegment.end = getDateFromPath(file); - - if (isLive(radio, file)) { - //log.info("listenHandler: currentAudioFile"); - var pData = { class: list[file].class, title: getRadio(radio).liveStatus.metadata }; - if (!isSameSegment(curSegment, pData)) { - result.push(curSegment); - curSegment = newSegment(index); - } - curSegment.title = pData.title; - filesCallback(); - } else { - fs.readFile(file + ".json", function(err, data) { - if (err) log.error("metadata: readFile err=" + err); - try { - var pData = JSON.parse(data); - } catch(e) { - log.error("metadata: json parse error file=" + file + ".json"); - return filesCallback(); - } - pData.class = list[file].class; - if (!isSameSegment(curSegment, pData)) { - result.push(curSegment); - curSegment = newSegment(index); - } - curSegment.title = pData.title; - filesCallback(); - }); - } - - - }, function(err) { - if (err) log.error("metadata: err=" + err.message); - result.push(curSegment); - response.set({ 'Access-Control-Allow-Origin': '*' }); - response.json(result); - }); - }); -} - - -app.get('/:action/:radio/:after', function(request, response) { +app.get('/:action/:radio/:delay', function(request, response) { var action = request.params.action; var radio = request.params.radio; - var after = request.params.after; - log.debug("get: action=" + action + " radio=" + radio + " after=" + after); + var delay = request.params.delay; + //log.debug("get: action=" + action + " radio=" + radio + " delay=" + delay); + if (!getRadio(radio) || !getRadio(radio).enable) { + response.writeHead(400); + return response.end(); + } switch(action) { case "listen": var ext = ".mp3"; // TODO check for other extensions - var state = { + /*var state = { nSegmentsInitialBuffer: Math.floor(LISTEN_BUFFER/SEG_DURATION), lastSentDate: new Date(), ext: ext, requestDate: new Date() } - listenRequestDate = state.requestDate; + listenRequestDate = state.requestDate;*/ - switch(state.ext) { + switch(ext) { case ".aac": response.set('Content-Type', 'audio/aacp'); break; case ".mp3": response.set('Content-Type', 'audio/mpeg'); break; } - listenHandler(response, radio, after, state, function() { + /*listenHandler(response, radio, after, state, function() { + response.end(); + });*/ + listenHandler(response, radio, delay, { + newRequest: true, + requestDate: new Date() + }, function() { response.end(); }); break; case "metadata": - metadataHandler(response, radio, after); + //metadataHandler(response, radio); + response.set({ 'Access-Control-Allow-Origin': '*' }); + var result = getRadio(radio).liveStatus.metaCache.read(); + response.json(result); break; default: - response.setHeader(400); + response.writeHead(400); response.end(); } });