diff --git a/lib/agent.js b/lib/agent.js index a7bdf6f..26b4ac2 100644 --- a/lib/agent.js +++ b/lib/agent.js @@ -88,7 +88,7 @@ var start = function start(done) { app._backlog = 2048; apps.push(app); - var secureOptions; + var secureOptions, options; dirModule = path.dirname(module.filename); logger.info('config.enableSecure', config.enableSecure); @@ -114,7 +114,7 @@ var start = function start(done) { fs.statSync(secureOptions.key).isFile() && fs.statSync(secureOptions.cert).isFile()) { - var options = { + options = { key: fs.readFileSync(secureOptions.key), cert: fs.readFileSync(secureOptions.cert) }; @@ -128,7 +128,7 @@ var start = function start(done) { appSec.secOptions = options; appSec.prefix = 'SEC:'; appSec.isSecure = true; - appSec.port = Number(config.agent.port) + 1; + appSec.port = Number(config.agent.httpsPort); apps.push(appSec); } diff --git a/lib/baseConfig.js b/lib/baseConfig.js index e6f8f37..39c9bc9 100644 --- a/lib/baseConfig.js +++ b/lib/baseConfig.js @@ -34,7 +34,7 @@ exports.logger = {}; exports.logger.logLevel = 'debug'; exports.logger.inspectDepth = 1; exports.logger.Console = { - level: 'debug', timestamp: true + level: 'info', timestamp: true }; exports.logger.File = { level: 'debug', filename: dir_prefix + @@ -179,6 +179,7 @@ exports.agent.maxMessages = 1000; * @type {Number} */ exports.agent.port = 3001; +exports.agent.httpsPort = 3002; /** * Provision timeout diff --git a/lib/dataSrv.js b/lib/dataSrv.js index 48ec8d9..05df66e 100644 --- a/lib/dataSrv.js +++ b/lib/dataSrv.js @@ -397,57 +397,48 @@ var peek = function(appPrefix, queue, maxElems, callback) { var queueId = queue.id, fullQueueIdH = config.dbKeyQueuePrefix + 'H:' + appPrefix + queue.id, fullQueueIdL = config.dbKeyQueuePrefix + 'L:' + appPrefix + queue.id, - restElems = 0; + restElems = 0, + db = dbCluster.getDb(queueId); - dbCluster.getOwnDb(queueId, function(err, db) { - if (err) { - manageError(err, callback); - } else { - peekAux(db); - } - }); - - function peekAux(db) { - db.lrange(fullQueueIdH, 0, maxElems - 1, function onRangeH(errH, dataH) { - - var dataHlength = dataH.length; - if (errH) {//errH - manageError(errH, callback); + db.lrange(fullQueueIdH, 0, maxElems - 1, function onRangeH(errH, dataH) { - } else { - if (dataHlength < maxElems) { + var dataHlength = dataH.length; + if (errH) {//errH + manageError(errH, callback); - restElems = maxElems - dataHlength; - //Extract from both queues - db.lrange(fullQueueIdL, 0, restElems - 1, - function on_rangeL(errL, dataL) { + } else { + if (dataHlength < maxElems) { - if (errL) { + restElems = maxElems - dataHlength; + //Extract from both queues + db.lrange(fullQueueIdL, 0, restElems - 1, + function on_rangeL(errL, dataL) { - //fail but we may have data of previous range - if (dataH) { - //if there is dataH dismiss the low priority error - getPeekData(dataH, callback, queue); - } else { - manageError(errL, callback); - } + if (errL) { + //fail but we may have data of previous range + if (dataH) { + //if there is dataH dismiss the low priority error + getPeekData(dataH, callback, queue); } else { + manageError(errL, callback); + } - if (dataL) { - dataH = dataH.concat(dataL); - } + } else { - getPeekData(dataH, callback, queue); + if (dataL) { + dataH = dataH.concat(dataL); } - }); - } else { - //just one queue used - getPeekData(dataH, callback, queue); - } + + getPeekData(dataH, callback, queue); + } + }); + } else { + //just one queue used + getPeekData(dataH, callback, queue); } - }); - } + } + }); }; function getPeekData(dataH, callback, queue) { @@ -758,14 +749,10 @@ var repushUndeliveredTransaction = function(appPrefix, queue, priority, extTrans var setBlockedQueue = function(appPrefix, queueId, blocked, cb) { - dbCluster.getOwnDb(queueId, function(err, db) { - if (err) { - manageError(err, cb); - } else { - var inc = blocked ? 1 : -1; - db.hincrby(config.dbKeyQueuePrefix + appPrefix + queueId + ':userInfo', 'blocked', inc, cb); - } - }); + var inc = blocked ? 1 : -1, + db = dbCluster.getDb(queueId); + + db.hincrby(config.dbKeyQueuePrefix + appPrefix + queueId + ':userInfo', 'blocked', inc, cb); }; var ackTransaction = function ackTransaction(extTransactionId, queue, cb) { diff --git a/package.json b/package.json index 029f444..e72df07 100644 --- a/package.json +++ b/package.json @@ -28,7 +28,9 @@ "hooker":"*", "logger":"https://github.com/telefonicaid/pditclogger/tarball/master", "performanceFramework":"https://github.com/PDI-DGS-Protolab/performanceFramework/tarball/master", - "should": "1.2.2", + "should": "1.2.2", + "forever": "0.10.8", + "daemon": "1.1.0", "socket.io": "0.8.7", "socket.io-client": "0.9.11" }, diff --git a/rpm/SOURCES/etc/init.d/pdi-popbox b/rpm/SOURCES/etc/init.d/pdi-popbox new file mode 100755 index 0000000..e4d76ef --- /dev/null +++ b/rpm/SOURCES/etc/init.d/pdi-popbox @@ -0,0 +1,99 @@ +#!/bin/bash +# description: Script to wake up Agent.js from Popbox module +# chkconfig: 2345 85 11 +########################################### +##Author: Carlos Enrique Gómez Gómez +##Depar: Release Engineering & Management +########################################### +## Values +# processname: popbox +# config: /opt/pdi-popbox/lib/baseConfig.js +# pidfile: /var/run/pdi-popbox/popbox.pid + +source /etc/rc.d/init.d/functions + +NAME=popbox +SOURCE_DIR=/opt/pdi-popbox/bin +SOURCE_FILE=popbox + +user=popbox +pidfile=/var/run/pdi-popbox/$NAME.pid +logfile=/var/log/pdi-popbox/$NAME.log +forever_dir=/var/run/forever + +node=node +sed=sed +home_dir="/opt/pdi-popbox" +forever=$home_dir/node_modules/forever/bin/forever + +start() { + echo "Starting $NAME node instance: " + + if [ "$foreverid" == "" ]; then + # Create the log and pid files, making sure that + # the target use has access to them + mkdir -p $(dirname $logfile) + touch $logfile + chown -R $user $(dirname $logfile) + mkdir -p $(dirname $pidfile) + #touch $pidfile + chown $user $(dirname $pidfile) + + # Launch the application + cd $home_dir + daemon --user=$user \ + $forever start -p $forever_dir --pidFile $pidfile -l $logfile -a \ + $SOURCE_DIR/$SOURCE_FILE + RETVAL=$? + echo + else + echo "Instance already running" + RETVAL=0 + echo + fi +} + +stop() { + [ "$pid" == "" ] && return ; + echo "Shutting down $NAME(pid:$pid) node instance : " + processes=$(ps -ef|grep ^$user|grep $pid|grep -v "$home_dir/node_modules/forever"|awk '{print $2" "$3}' 2>/dev/null) + [ "$processes" == "" ] && echo "Any $user proccess is up" && RETVAL=1 + if [ "$RETVAL" != "1" ]; then + kill -9 $processes + RETVAL=$? + fi + rm -Rf $pidfile + [ "$RETVAL" == "0" ] && echo_success + echo +} + +if [ -f $pidfile ]; then + read pid < $pidfile +else + pid="" +fi +if [ "$pid" != "" ]; then + # Gnarly sed usage to obtain the foreverid. + sed1="/$pid\]/p" + sed2="s/.*\[\([0-9]\+\)\].*\s$pid\].*/\1/g" + foreverid=`$forever list -p $forever_dir | $sed -n $sed1 | $sed $sed2` +else + foreverid="" +fi + +case "$1" in + start) + start + ;; + stop) + stop + ;; + status) + status -p ${pidfile} + ;; + *) + echo "Usage: {start|stop|status}" + exit 1 + ;; +esac +exit $RETVAL diff --git a/rpm/SOURCES/etc/logrotate.d/pdi-popbox b/rpm/SOURCES/etc/logrotate.d/pdi-popbox new file mode 100644 index 0000000..bf4364a --- /dev/null +++ b/rpm/SOURCES/etc/logrotate.d/pdi-popbox @@ -0,0 +1,9 @@ +/var/log/popbox/*.log{ + daily + rotate 10 + copytruncate + delaycompress + compress + notifempty + missingok +} \ No newline at end of file diff --git a/rpm/SPECS/popbox.spec b/rpm/SPECS/popbox.spec new file mode 100644 index 0000000..efb8cbf --- /dev/null +++ b/rpm/SPECS/popbox.spec @@ -0,0 +1,112 @@ +Summary: Popbox module to manage node + Redis value server +Name: popbox +Version: 0.0.2 +Release: 2 +License: GNU +BuildRoot: %{_topdir}/BUILDROOT/ +BuildArch: x86_64 +Requires: nodejs >= 0.8 +Requires(post): /sbin/chkconfig /usr/sbin/useradd +Requires(preun): /sbin/chkconfig, /sbin/service +Requires(postun): /sbin/service +Group: Applications/Popbox +Vendor: Telefonica I+D +BuildRequires: npm + +%description +Simple High-Performance High-Scalability Inbox Notification Service, +require Redis 2.6 server, node 0.8 and npm only for installatión +%define _prefix_company pdi- +%define _project_name popbox +%define _project_user %{_project_name} +%define _company_project_name %{_prefix_company}%{_project_name} +%define _service_name %{_company_project_name} +%define _install_dir /opt +%define _srcdir %{_sourcedir}/../../ +%define _project_install_dir %{_install_dir}/%{_company_project_name} +%define _logrotate_conf_dir %{_srcdir}/conf/log +%define _popbox_log_dir %{_localstatedir}/log/%{_company_project_name} +%define _build_root_project %{buildroot}/%{_project_install_dir} +%define _conf_dir /etc/%{_prefix_company}%{_project_name} +%define _log_dir /var/log/%{_prefix_company}%{_project_name} + + +# -------------------------------------------------------------------------------------------- # +# prep section, setup macro: +# -------------------------------------------------------------------------------------------- # +%prep +rm -Rf $RPM_BUILD_ROOT && mkdir -p $RPM_BUILD_ROOT +[ -d %{_build_root_project} ] || mkdir -p %{_build_root_project} + +cp -R %{_srcdir}/lib %{_srcdir}/package.json %{_srcdir}/index.js \ + %{_srcdir}/bin %{_srcdir}/License.txt %{_build_root_project} +cp -R %{_sourcedir}/* %{buildroot} +mkdir -p %{buildroot}/var/run/%{_company_project_name} +mkdir -p %{buildroot}/var/log/%{_company_project_name} + +%build +cd %{_build_root_project} +# Only production modules +npm install --production +rm package.json + +# -------------------------------------------------------------------------------------------- # +# pre-install section: +# -------------------------------------------------------------------------------------------- # +%pre +echo "[INFO] Creating %{_project_user} user" +grep ^%{_project_user} /etc/passwd +RET_VAL=$? +if [ "$RET_VAL" != "0" ]; then + /usr/sbin/useradd -c '%{_project_user}' -u 699 -s /bin/false \ + -r -d %{_project_install_dir} %{_project_user} + RET_VAL=$? + if [ "$RET_VAL" != "0" ]; then + echo "[ERROR] Unable create popbox user" \ + exit $RET_VAL + fi +fi + +# -------------------------------------------------------------------------------------------- # +# post-install section: +# -------------------------------------------------------------------------------------------- # +%post +echo "Configuring application:" +rm -Rf /etc/initi.d/%{_service_name} +cd /etc/init.d + +#Service +echo "Creating %{_service_name} service:" +chkconfig --add %{_service_name} + +#Config +rm -Rf %{_conf_dir} && mkdir -p %{_conf_dir} +cd %{_conf_dir} +ln -s %{_project_install_dir}/lib/baseConfig.js %{_project_name}_conf.js + +#Logs +#TODO configuration logs +echo "Done" + +%preun +if [ $1 == 0 ]; then + echo "Removing application config files" + [ -d /etc/%{_company_project_name} ] && rm -rfv /etc/%{_company_project_name} + [ -d %{_popbox_log_dir} ] && rm -rfv %{_popbox_log_dir} + [ -d %{_project_install_dir} ] && rm -rfv %{_project_install_dir} + echo "Destroying %{_service_name} service:" + chkconfig --del %{_service_name} + rm -Rf /etc/init.d/%{_service_name} + + echo "Done" +fi + +%postun +%clean +rm -rf $RPM_BUILD_ROOT +%files +%defattr(755,%{_project_user},%{_project_user},755) +%config /etc/init.d/%{_service_name} +%config /etc/logrotate.d/%{_company_project_name} +%{_project_install_dir} +/var/ \ No newline at end of file diff --git a/rpm/package.sh b/rpm/package.sh new file mode 100755 index 0000000..96473f4 --- /dev/null +++ b/rpm/package.sh @@ -0,0 +1,3 @@ +#!/bin/bash +DIR=$(dirname $(readlink -f $0)) +rpmbuild -v --clean -ba $DIR/SPECS/popbox.spec --define '_topdir '$DIR diff --git a/utils/ack.js b/utils/ack.js new file mode 100644 index 0000000..2cdf1ee --- /dev/null +++ b/utils/ack.js @@ -0,0 +1,42 @@ +#!/usr/bin/env node + +'use strict'; + +var program = require('commander'), + request = require('request'), + util = require('util'), + i, + url, + trans; + +program + .version('0.0.1') + .option('-H, --host [hostname]', 'host, \'localhost\' by default', 'localhost') + .option('-P, --port [number]', 'port, 5001 by default', 3001, parseInt) + .option('-Q, --queue [id]', 'queue', "Q1") + .option('-T, --trans [list]', 'list of transactions separated by comma', list) + .option('-X, --secure', 'use HTTPS', false) + .parse(process.argv); +url = util.format('%s://%s:%d/queue/%s/ack', program.secure?'https':'http', program.host, program.port, program.queue); +trans = { + 'transactions': program.trans +}; + +console.log('url\n', url); +console.log('trans\n', trans); +request.post({url : url, json: trans}, function(err, res, body) { + if(err) { + console.log('error\n', err); + } + else { + console.log('statusCode\n', res.statusCode); + console.log('headers\n', res.headers); + console.log('body\n', body); + } +}); +function range(val) { + return val.split('..').map(Number); +} +function list(val) { + return val.split(','); +} diff --git a/utils/createSecure.js b/utils/createSecure.js index 1c763af..b7f53f8 100755 --- a/utils/createSecure.js +++ b/utils/createSecure.js @@ -12,7 +12,7 @@ var program = require('commander'), program .version('0.0.1') .option('-H, --host [hostname]', 'host, \'localhost\' by default', 'localhost') - .option('-P, --port [number]', 'port, 5001 by default', 5002, parseInt) + .option('-P, --port [number]', 'port, 5001 by default', 3002, parseInt) .option('-Q, --queue [id]', 'queue', "Q1") .option('-U, --user [id:passwd]', 'username','popbox:itscool') .parse(process.argv); diff --git a/utils/group.js b/utils/group.js index 61ab161..d69b578 100755 --- a/utils/group.js +++ b/utils/group.js @@ -12,7 +12,7 @@ var program = require('commander'), program .version('0.0.1') .option('-H, --host [hostname]', 'host, \'localhost\' by default', 'localhost') - .option('-P, --port [number]', 'port, 5001 by default', 5001, parseInt) + .option('-P, --port [number]', 'port, 5001 by default', 3001, parseInt) .option('-Q, --queues [list]', 'list of queues separated by comma', list, ["Q1"]) .option('-G, --group [id]', 'group to be created', ["G1"]) .option('-X, --secure', 'use HTTPS', false) diff --git a/utils/pop.js b/utils/pop.js index 292fead..c975968 100755 --- a/utils/pop.js +++ b/utils/pop.js @@ -17,7 +17,7 @@ var program = require('commander'), program .version('0.0.1') .option('-H, --host [hostname]', 'host, \'localhost\' by default', 'localhost') - .option('-P, --port [number]', 'port, 5001 by default', 5001, parseInt) + .option('-P, --port [number]', 'port, 5001 by default', 3001, parseInt) .option('-Q, --queue [id]', 'queue', "Q1") .option('-R, --reliable', 'reliable extraction', false) .option('-S, --subscribe', 'subscribe', false) diff --git a/utils/provision.js b/utils/provision.js index 7fa0932..b1ed62e 100755 --- a/utils/provision.js +++ b/utils/provision.js @@ -12,7 +12,7 @@ var program = require('commander'), program .version('0.0.1') .option('-H, --host [hostname]', 'host, \'localhost\' by default', 'localhost') - .option('-P, --port [number]', 'port, 5001 by default', 5001, parseInt) + .option('-P, --port [number]', 'port, 5001 by default', 3001, parseInt) .option('-M, --message [text]', 'message','¡hola!\n') .option('-Q, --queues [list]', 'list of queues separated by comma', list, ["Q1"]) .option('-G, --groups [list]', 'list of groups separated by comma',list) @@ -28,8 +28,8 @@ trans = { if(program.expiration) { trans.expirationDelay = program.expiration; } -if(program.groups && program.groups.length!==0) { - trans.groups = program.groups; +if(program.groups) { + trans.groups = program.groups; } if(program.callback) { trans.callback = program.callback; diff --git a/utils/trans_info.js b/utils/trans_info.js index 5a21e54..b3c3de3 100755 --- a/utils/trans_info.js +++ b/utils/trans_info.js @@ -12,15 +12,15 @@ var program = require('commander'), program .version('0.0.1') .option('-H, --host [hostname]', 'host, \'localhost\' by default', 'localhost') - .option('-P, --port [number]', 'port, 5001 by default', 5001, parseInt) + .option('-P, --port [number]', 'port, 5001 by default', 3001, parseInt) .option('-T, --trans [id]', 'transaction') .option('-X, --secure', 'use HTTPS', false) .parse(process.argv); url = util.format('%s://%s:%d/trans/%s', program.secure ? 'https' : 'http', program.host, program.port, program.trans); -req_aux(url, function () { - req_aux(url + '/state', function(){}); +reqAux(url, function () { + reqAux(url + '/state', function(){}); }); -function req_aux(url, done) { +function reqAux(url, done) { console.log('url\n', url); request.get({url: url, headers: {'Accept': 'application/json'}}, function (err, res, body) { if (err) {