server works with delay for audio, delay progressbar in UI

This commit is contained in:
Alexandre Storelli
2018-01-18 00:12:09 +01:00
parent 5bae1ebbad
commit 4f5fb590bf
10 changed files with 742 additions and 274 deletions
+160 -10
View File
@@ -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;
+2 -1
View File
@@ -17,5 +17,6 @@
"build": "react-scripts build",
"test": "react-scripts test --env=jsdom",
"eject": "react-scripts eject"
}
},
"homepage" : "https://www.adblockradio.com/buffer"
}
+52 -40
View File
@@ -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 = (
<span>
{self.state.playingRadio.split("_")[1]}
{self.state.playingLive ?
" (direct)"
:
" (" + moment(+self.state.date - self.state.playingDelta + LIVE_DELAY).fromNow() + ")"
}
{delayText}
</span>
)
} else {
@@ -177,6 +184,7 @@ class App extends Component {
</StatusButtonsContainer>
);
//console.log("Metadata props: date=" + (+self.state.date) + " clockDiff=" + self.state.clockDiff + " playingDelay=" + self.state.playingDelay);
return (
<AppParent>
@@ -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} />
)
})}
</RadioList>
{self.state.playingRadio &&
<Metadata playingRadio={self.state.playingRadio}
metaList={self.state["meta" + self.state.playingRadio]}
playingDate={+self.state.date - self.state.playingDelta}
currentDate={self.state.date}
playingDelay={self.state.playingDelay}
metaList={self.state[self.state.playingRadio + "|metadata"]}
date={+self.state.date - self.state.clockDiff}
playCallback={self.play} />
}
</AppView>
<Controls>
{status}
{buttons}
<DelayCanvas playing={!!self.state.playingRadio}
playingDelay={self.state.playingDelay}
cacheLen={self.state.config.user.cacheLen}
playCallback={self.play} />
{/*<PlayerStatus settings={this.props.settings} bsw={this.props.bsw} condensed={this.props.condensed} playbackAction={this.togglePlayer} />*/}
</Controls>
</AppParent>
@@ -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`
+88
View File
@@ -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 (
<canvas width={canvasStyle.width} height={canvasStyle.height} ref="canvas" id="canvas" style={canvasStyle} />
)
}
}
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;
+8 -7
View File
@@ -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 {
<MetadataContainer className={classNames({ compact: previewMode })}>
{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 (
<MetadataItem className={classNames({ playing: playing, compact: previewMode })} key={"item" + i} onClick={function() { self.play(item.start); }}>
<MetadataItem className={classNames({ playing: playing, compact: previewMode })} key={"item" + i} onClick={function() { self.play(item.validFrom); }}>
<MetadataText>
{!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}
</MetadataText>
<MetadataCover className={classNames({ playing: playing })} src={item.title.cover || defaultCover} alt="logo" />
<MetadataCover className={classNames({ playing: playing })} src={item.payload.cover || defaultCover} alt="logo" />
</MetadataItem>
)
})
@@ -79,7 +80,7 @@ const MetadataItem = styled.div`
background: #eee;
display: flex;
cursor: pointer;
&.playing {
border: 2px solid red;
}
+3 -2
View File
@@ -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);
+15
View File
@@ -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;
+139
View File
@@ -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<consts.WLARRAY.length; i++) {
files[consts.WLARRAY[i]] = {};
}
getDirs(options.path + "/records", function(dateDirs) {
for (let i=dateDirs.length-1; 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;
+234
View File
@@ -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<config.radios.length; j++) {
if (config.radios[j].country == country && config.radios[j].name == name) {
return config.radios[j];
}
}
} else { // only first parameter used
for (var j=0; j<config.radios.length; j++) {
if (config.radios[j].country + "_" + config.radios[j].name == country) {
return config.radios[j];
}
}
}
return null;
}
exports.getRadio = getRadio;
var isLive = function(radio, file) {
return getRadio(radio).liveStatus.currentPrefix == file;
}
/*var listenRequestDate = null;
exports.listenHandler = function(response, radio, after, state, callback) {
var before = new Date(+new Date(after)+6*SEG_DURATION*1000).toISOString();
log.debug("listenHandler: after=" + after + " before=" + before + " bufferSeg=" + state.nSegmentsInitialBuffer);
findDataFiles({ radios: [ radio ], after: after, before: before, path: __dirname }, function(classes) {
var list = {};
for (var classItem in classes) {
Object.assign(list, classes[classItem]);
}
var files = Object.keys(list);
files.sort();
if (files.length == 0) return callback();
var timeoutMonitor = null;
async.forEachOfSeries(files, function(file, index, filesCallback) {
log.debug("listen: read=" + file + state.ext);
var rs;
var livePlay = isLive(radio, file);
if (livePlay) {
rs = getRadio(radio).liveStatus.liveReadStream;
} else {
rs = fs.createReadStream(file + state.ext);
}
var dataRead = 0;
rs.on("data", function(data) {
response.write(data);
dataRead += data.length;
});
rs.on("error", function(err) {
log.error("listen: readFile err=" + err + " livePlay=" + livePlay);
});
rs.on("end", function() {
state.lastSentDate = new Date();
var willWaitDrain = !response.write("");
log.debug("listen: sent file " + (index+1) + "/" + files.length + " bytes=" + dataRead + " waitDrain=" + willWaitDrain + " live=" + livePlay);
if (livePlay || !willWaitDrain) { //willWaitDrain) { // detect congestion of stream
bufManager();
} else {
log.debug("listenHandler: will wait for drain event");
response.once("drain", function() {
clearTimeout(timeoutMonitor);
bufManager();
});
timeoutMonitor = setTimeout(function() {
response.removeListener("drain", bufManager);
response.destroy();
filesCallback({ message: "drain event not emitted, connection timeout" });
}, SEG_DURATION*1500);
}
});
var bufManager = function() {
before = new Date(+new Date(getDateFromPath(file))+1).toISOString();
if (state.requestDate !== listenRequestDate) {
filesCallback({ message: "request canceled because another one has been initiated" });
} else if (livePlay) {
filesCallback();
} else if (state.nSegmentsInitialBuffer > 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);
});
});
}*/
+41 -214
View File
@@ -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<config.radios.length; i++) { // control on what data is exposed via the api
var radio = config.radios[i];
if (!radio.enable) continue;
result.radios.push({
country: radio.country,
name: radio.name,
enable: radio.enable,
content: radio.content,
url: radio.enabled,
favicon: radio.favicon
});
}
response.json(result);
//log.debug("config[0]=" + JSON.stringify(config.radios[0]));
});
var getDateFromPath = function(path) {
var spl = path.split("/");
return spl[spl.length-1];
}
var { listenHandler, metadataHandler, getRadio } = require("./handlers.js");
var getRadio = function(country, name) {
if (name) { // both parameters used
for (var j=0; j<config.radios.length; j++) {
if (config.radios[j].country == country && config.radios[j].name == name) {
return config.radios[j];
}
}
} else { // only first parameter used
for (var j=0; j<config.radios.length; j++) {
if (config.radios[j].country + "_" + config.radios[j].name == country) {
return config.radios[j];
}
}
}
return null;
}
var isLive = function(radio, file) {
return getRadio(radio).liveStatus.currentPrefix == file;
}
var listenRequestDate = null;
var listenHandler = function(response, radio, after, state, callback) {
var before = new Date(+new Date(after)+6*SEG_DURATION*1000).toISOString();
log.debug("listenHandler: after=" + after + " before=" + before + " bufferSeg=" + state.nSegmentsInitialBuffer);
findDataFiles({ radios: [ radio ], after: after, before: before, path: __dirname }, function(classes) {
var list = {};
for (var classItem in classes) {
Object.assign(list, classes[classItem]);
}
var files = Object.keys(list);
files.sort();
if (files.length == 0) return callback();
var timeoutMonitor = null;
async.forEachOfSeries(files, function(file, index, filesCallback) {
log.debug("listen: read=" + file + state.ext);
var rs;
var livePlay;
if (isLive(radio, file)) {
livePlay = true;
//log.info("listenHandler: currentAudioFile");
rs = getRadio(radio).liveStatus.liveReadStream;
} else {
livePlay = false;
rs = fs.createReadStream(file + state.ext);
}
var dataRead = 0;
rs.on("data", function(data) {
response.write(data);
dataRead += data.length;
});
rs.on("error", function(err) {
log.error("listen: readFile err=" + err + " livePlay=" + livePlay);
});
rs.on("end", function() {
state.lastSentDate = new Date();
var willWaitDrain = !response.write("");
log.debug("listen: sent file " + (index+1) + "/" + files.length + " bytes=" + dataRead + " waitDrain=" + willWaitDrain + " live=" + livePlay);
if (livePlay || !willWaitDrain) { //willWaitDrain) { // detect congestion of stream
bufManager();
} else {
log.debug("listenHandler: will wait for drain event");
response.once("drain", function() {
clearTimeout(timeoutMonitor);
bufManager();
});
timeoutMonitor = setTimeout(function() {
response.removeListener("drain", bufManager);
response.destroy();
filesCallback({ message: "drain event not emitted, connection timeout" });
}, SEG_DURATION*1500);
}
});
var bufManager = function() {
before = new Date(+new Date(getDateFromPath(file))+1).toISOString();
if (state.requestDate !== listenRequestDate) {
filesCallback({ message: "request canceled because another one has been initiated" });
} else if (livePlay) {
filesCallback();
} else if (state.nSegmentsInitialBuffer > 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();
}
});