CC-2588: Use RabbitMQ to control the show recorder

- recording works now.
- further testing is needed
- still need to work on canceling the show
This commit is contained in:
james 2011-07-23 12:30:36 -04:00
parent ccdb886b9a
commit 047a288c38
4 changed files with 273 additions and 89 deletions

View File

@ -37,6 +37,8 @@ class RabbitMq
$channel->basic_publish($msg, $EXCHANGE); $channel->basic_publish($msg, $EXCHANGE);
$channel->close(); $channel->close();
$conn->close(); $conn->close();
self::SendMessageToShowRecorder("update_schedule");
} }
} }
@ -63,5 +65,26 @@ class RabbitMq
$channel->close(); $channel->close();
$conn->close(); $conn->close();
} }
public static function SendMessageToShowRecorder($event_type)
{
global $CC_CONFIG;
$conn = new AMQPConnection($CC_CONFIG["rabbitmq"]["host"],
$CC_CONFIG["rabbitmq"]["port"],
$CC_CONFIG["rabbitmq"]["user"],
$CC_CONFIG["rabbitmq"]["password"]);
$channel = $conn->channel();
$channel->access_request($CC_CONFIG["rabbitmq"]["vhost"], false, false, true, true);
$EXCHANGE = 'airtime-show-recorder';
$channel->exchange_declare($EXCHANGE, 'direct', false, true);
$msg = new AMQPMessage($event_type, array('content_type' => 'text/plain'));
$channel->basic_publish($msg, $EXCHANGE);
$channel->close();
$conn->close();
}
} }

View File

@ -0,0 +1,144 @@
#!/usr/local/bin/python
import urllib
import logging
import logging.config
import json
import time
import datetime
import os
import sys
import shutil
from configobj import ConfigObj
from poster.encode import multipart_encode
from poster.streaminghttp import register_openers
import urllib2
from subprocess import Popen
from threading import Thread
import mutagen
from api_clients import api_client
# For RabbitMQ - to be implemented in the future
from kombu.connection import BrokerConnection
from kombu.messaging import Exchange, Queue, Consumer, Producer
# configure logging
try:
logging.config.fileConfig("logging.cfg")
except Exception, e:
print 'Error configuring logging: ', e
sys.exit()
POLL_INTERVAL=3600
# loading config file
try:
config = ConfigObj('/etc/airtime/recorder.cfg')
except Exception, e:
logger = logging.getLogger('root')
logger.error('Error loading config file: %s', e)
sys.exit()
def getDateTimeObj(time):
timeinfo = time.split(" ")
date = timeinfo[0].split("-")
time = timeinfo[1].split(":")
return datetime.datetime(int(date[0]), int(date[1]), int(date[2]), int(time[0]), int(time[1]), int(time[2]))
class CommandHandler(Thread):
def __init__(self, command, q):
Thread.__init__(self)
self.api_client = api_client.api_client_factory(config)
self.logger = logging.getLogger('root')
self.logger.info("Handling command: " + command)
self.command = command
self.queue = q
def run(self):
if(self.command == 'update_schedule'):
self.queue.put(self.api_client.get_shows_to_record())
class CommandListener(Thread):
def __init__(self, q):
Thread.__init__(self)
self.api_client = api_client.api_client_factory(config)
self.logger = logging.getLogger('root')
self.sr = None
self.queue = q
self.shows_to_record = []
self.logger.info("RecorderFetch: init complete")
def init_rabbit_mq(self):
self.logger.info("Initializing RabbitMQ stuff")
try:
schedule_exchange = Exchange("airtime-show-recorder", "direct", durable=True, auto_delete=True)
schedule_queue = Queue("recorder-fetch", exchange=schedule_exchange, key="foo")
self.connection = BrokerConnection(config["rabbitmq_host"], config["rabbitmq_user"], config["rabbitmq_password"], "/")
channel = self.connection.channel()
consumer = Consumer(channel, schedule_queue)
consumer.register_callback(self.handle_message)
consumer.consume()
except Exception, e:
self.logger.error(e)
return False
return True
def handle_message(self, body, message):
self.logger.info("Received command from RabbitMQ: " + message.body)
ch = CommandHandler(message.body, self.queue)
ch.start()
# ACK the message to take it off the queue
message.ack()
"""def process_shows(self, shows):
self.logger.info("Processing shows...")
self.shows_to_record = []
self.logger.info(shows)
shows = shows[u'shows']
for show in shows:
show_starts = getDateTimeObj(show[u'starts'])
show_end = getDateTimeObj(show[u'ends'])
time_delta = show_end - show_starts
self.shows_to_record[show[u'starts']] = [time_delta, show[u'instance_id'], show[u'name']]"""
"""
Main loop of the thread:
Wait for schedule updates from RabbitMQ, but in case there arent any,
poll the server to get the upcoming schedule.
"""
def run(self):
self.logger.info("Started...")
while not self.init_rabbit_mq():
self.logger.error("Error connecting to RabbitMQ Server. Trying again in few seconds")
time.sleep(5)
# Bootstrap: since we are just starting up, we need to grab the
# most recent schedule. After that we can just wait for updates.
try:
self.queue.put(self.api_client.get_shows_to_record())
self.logger.info("Bootstrap complete: got initial copy of the schedule")
except Exception, e:
self.logger.error(e)
loops = 1
while True:
self.logger.info("Loop #%s", loops)
try:
# Wait for messages from RabbitMQ. Timeout if we
# dont get any after POLL_INTERVAL.
self.connection.drain_events(timeout=POLL_INTERVAL)
status = 1
except Exception, e:
self.logger.info(e)
# We didnt get a message for a while, so poll the server
# to get an updated schedule.
self.queue.put(self.api_client.get_shows_to_record())
loops += 1

View File

@ -8,3 +8,11 @@ base_recorded_files = '/var/tmp/airtime/show-recorder/'
# where the logging files live # where the logging files live
log_dir = '/var/log/airtime/show-recorder' log_dir = '/var/log/airtime/show-recorder'
############################################
# RabbitMQ settings #
############################################
rabbitmq_host = 'localhost'
rabbitmq_user = 'guest'
rabbitmq_password = 'guest'

View File

@ -8,6 +8,9 @@ import datetime
import os import os
import sys import sys
import shutil import shutil
from Queue import Queue
from commandlistener import CommandListener
from configobj import ConfigObj from configobj import ConfigObj
@ -20,10 +23,6 @@ from threading import Thread
import mutagen import mutagen
# For RabbitMQ - to be implemented in the future
#from kombu.connection import BrokerConnection
#from kombu.messaging import Exchange, Queue, Consumer, Producer
from api_clients import api_client from api_clients import api_client
# configure logging # configure logging
@ -61,7 +60,7 @@ class ShowRecorder(Thread):
self.show_name = show_name self.show_name = show_name
self.logger = logging.getLogger('root') self.logger = logging.getLogger('root')
self.p = None self.p = None
def record_show(self): def record_show(self):
length = str(self.filelength)+".0" length = str(self.filelength)+".0"
filename = self.start_time filename = self.start_time
@ -114,109 +113,119 @@ class ShowRecorder(Thread):
datagen, headers = multipart_encode({"file": open(filepath, "rb"), 'name': filename, 'show_instance': self.show_instance}) datagen, headers = multipart_encode({"file": open(filepath, "rb"), 'name': filename, 'show_instance': self.show_instance})
self.api_client.upload_recorded_show(datagen, headers) self.api_client.upload_recorded_show(datagen, headers)
def set_metadata_and_save(self, filepath):
try:
date = self.start_time
md = date.split(" ")
time = md[1].replace(":", "-")
self.logger.info("time: %s" % time)
name = time+"-"+self.show_name
name.encode('utf-8')
artist = "AIRTIMERECORDERSOURCEFABRIC".encode('utf-8')
#set some metadata for our file daemon
recorded_file = mutagen.File(filepath, easy=True)
recorded_file['title'] = name
recorded_file['artist'] = artist
recorded_file['date'] = md[0]
recorded_file['tracknumber'] = self.show_instance
recorded_file.save()
except Exception, e:
self.logger.error("Exception: %s", e)
def run(self): def run(self):
code, filepath = self.record_show() code, filepath = self.record_show()
if code == 0: if code == 0:
self.logger.info("Preparing to upload %s" % filepath)
try: try:
date = self.start_time self.logger.info("Preparing to upload %s" % filepath)
md = date.split(" ")
time = md[1].replace(":", "-") self.set_metadata_and_save(filepath)
self.logger.info("time: %s" % time)
self.upload_file(filepath)
name = time+"-"+self.show_name except Exceptio, e:
name.encode('utf-8') self.logger.error(e)
artist = "AIRTIMERECORDERSOURCEFABRIC".encode('utf-8')
#set some metadata for our file daemon
recorded_file = mutagen.File(filepath, easy=True)
recorded_file['title'] = name
recorded_file['artist'] = artist
recorded_file['date'] = md[0]
recorded_file['tracknumber'] = self.show_instance
recorded_file.save()
except Exception, e:
self.logger.error("Exception: %s", e)
self.upload_file(filepath)
else: else:
self.logger.info("problem recording show") self.logger.info("problem recording show")
class RecordScheduler(Thread):
class Record(): def __init__(self, q):
Thread.__init__(self)
def __init__(self): self.queue = q
self.api_client = api_client.api_client_factory(config)
self.shows_to_record = {} self.shows_to_record = {}
self.logger = logging.getLogger('root') self.logger = logging.getLogger('root')
self.sr = None
def process_shows(self, shows): def process_shows(self, shows):
self.logger.info("Processing show schedules...")
self.shows_to_record = {} self.shows_to_record = {}
temp = shows[u'shows']
for show in shows: for show in temp:
show_starts = getDateTimeObj(show[u'starts']) show_starts = getDateTimeObj(show[u'starts'])
show_end = getDateTimeObj(show[u'ends']) show_end = getDateTimeObj(show[u'ends'])
time_delta = show_end - show_starts time_delta = show_end - show_starts
self.shows_to_record[show[u'starts']] = [time_delta, show[u'instance_id'], show[u'name']] self.shows_to_record[show[u'starts']] = [time_delta, show[u'instance_id'], show[u'name']]
self.logger.info(self.shows_to_record)
def check_record(self): def check_record(self):
if len(self.shows_to_record) != 0:
tnow = datetime.datetime.now() try:
sorted_show_keys = sorted(self.shows_to_record.keys()) tnow = datetime.datetime.now()
sorted_show_keys = sorted(self.shows_to_record.keys())
start_time = sorted_show_keys[0]
next_show = getDateTimeObj(start_time) start_time = sorted_show_keys[0]
next_show = getDateTimeObj(start_time)
self.logger.debug("Next show %s", next_show)
self.logger.debug("Now %s", tnow) self.logger.debug("Next show %s", next_show)
self.logger.debug("Now %s", tnow)
delta = next_show - tnow
min_delta = datetime.timedelta(seconds=60) delta = next_show - tnow
min_delta = datetime.timedelta(seconds=5)
if delta <= min_delta:
self.logger.debug("sleeping %s seconds until show", delta.seconds) if delta <= min_delta:
time.sleep(delta.seconds) self.logger.debug("sleeping %s seconds until show", delta.seconds)
time.sleep(delta.seconds)
show_length = self.shows_to_record[start_time][0]
show_instance = self.shows_to_record[start_time][1] show_length = self.shows_to_record[start_time][0]
show_name = self.shows_to_record[start_time][2] show_instance = self.shows_to_record[start_time][1]
show_name = self.shows_to_record[start_time][2]
self.sr = ShowRecorder(show_instance, show_name, show_length.seconds, start_time, filetype="mp3")
self.sr.start() self.sr = ShowRecorder(show_instance, show_name, show_length.seconds, start_time, filetype="mp3")
self.sr.start()
#remove show from shows to record.
del self.shows_to_record[start_time] #remove show from shows to record.
del self.shows_to_record[start_time]
except Exception,e :
def get_shows(self): self.logger.error(e)
else:
response = self.api_client.get_shows_to_record() self.logger.info("No recording schedule...")
if response is not None and 'is_recording' in response:
if self.sr is not None: def run(self):
if not response['is_recording'] and self.sr.is_recording(): self.logger.info("RecordScheduler started...")
self.sr.cancel_recording() while True:
if not self.queue.empty():
shows = response[u'shows'] try:
self.logger.debug("Received data from command handler")
if len(shows): shows = self.queue.get()
self.process_shows(shows) self.logger.debug('shows %s' % shows)
self.check_record() self.process_shows(shows)
except Exception, e:
self.logger.error(e)
self.check_record()
time.sleep(1)
if __name__ == '__main__': if __name__ == '__main__':
q = Queue()
recorder = Record()
cl = CommandListener(q)
while True: cl.start()
recorder.get_shows()
time.sleep(5) rs = RecordScheduler(q)
rs.start()