rendered paste body#!/usr/bin/python# -*- coding: utf-8 -*-import timeimport datetimeimport osimport urlparseimport statimport loggingimport pygstpygst.require("0.10")import gstimport gobjectimport dbusimport dbus.servicefrom dbus.mainloop.glib import DBusGMainLoopDBusGMainLoop(set_as_default=True)def get_element_from_url(url): res = urlparse.urlparse(url) if res.scheme == "http": return "souphttpsrc uri=%s"%url elif res.scheme == "rtsp": return "rtspsrc location=%s"%url elif res.scheme == "file" or not res.scheme: st = os.stat(res.path) if stat.S_ISCHR(st.st_mode): return "v4l2src device=%s"%res.path else: return "filesrc location=%s"%res.path return Nonedef get_monitor_arguments_from_string(s): def _get_bool(v): s = str(v).strip().lower() if s == "yes" or s == "true" or (s.isdigit() and float(s) != 0.0): return True elif s == "no" or s == "false" or (s.isdigit() and float(s) == 0.0): return False return None d = {} spl = s.split("!") for sp in spl: sp = sp.strip() spl2 = sp.split("=", 1) if len(spl2) == 2: key = spl2[0].strip() if key in d: raise ValueError("duplicate parameter '%s'"%key) value = spl2[1].strip() if key == "name": value = str(value) elif key == "source": source = str(value) elif key == "geometry": value = value.split("x") if len(value) != 2: raise ValueError("invalid geometry '%s'"%(value)) try: value = tuple((int(value[0].strip()), int(value[1].strip()))) except (ValueError, TypeError): raise ValueError("invalid geometry parameters '%s'"%(value)) elif key == "debug": value = _get_bool(value) if value is None: raise ValueError("invalid debug boolean '%s'"%value) elif key == "framerate": try: value = int(value) except (ValueError, TypeError): raise ValueError("invalid framerate '%s'"%(value)) if value < 1 or value > 256: raise ValueError("framerate out of range '%d'"%value) elif key == "alarm-pre-time": try: value = float(value) except (ValueError, TypeError): raise ValueError("invalid alarm-pre-time '%s'"%(value)) if value < 0.0 or value > 3600.0: raise ValueError("alarm-pre-time out of range '%f'"%value) elif key == "alarm-post-time": try: value = float(value) except (ValueError, TypeError): raise ValueError("invalid alarm-post-time '%s'"%(value)) if value < 0.0 or value > 36000.0: raise ValueError("alarm-post-time out of range '%f'"%value) else: raise ValueError("invalid parameter name '%s'"%key) d[key] = value elif not d: d["source"] = sp else: raise ValueError("invalid argument syntax '%s'"%sp) if "source" in d: source = d["source"] del d["source"] else: source = None if "name" in d: name = d["name"] del d["name"] else: name = None return name, source, dclass Monitor(gobject.GObject): STATE_INACTIVE = 0x00 STATE_IDLE = 0x10 STATE_PRE_ALARM = 0x30 STATE_ALARM = 0x40 STATE_POST_ALARM = 0x50 __gsignals__ = { "state-changed": (gobject.SIGNAL_RUN_FIRST, gobject.TYPE_NONE, (int, int)), } def __init__(self, name, source, **kwargs): gobject.GObject.__init__(self) self._name = name self._pipeline = None self._alarm_timestamp = 0.0 self._alarm_pre_time = float(kwargs.get("alarm-pre-time", 5.0)) self._alarm_post_time = float(kwargs.get("alarm-post-time", 5.0)) self._state = kwargs.get("state", Monitor.STATE_IDLE) if not self._state in (Monitor.STATE_IDLE, Monitor.STATE_PRE_ALARM, Monitor.STATE_ALARM, Monitor.STATE_POST_ALARM, ): raise ValueError("invalid state '%s'"%self._state) self.set_source(source) self.set_geometry(kwargs.get("geometry")) self.set_framerate(kwargs.get("framerate")) self.set_debug(kwargs.get("debug", False)) gobject.timeout_add_seconds(1, self.on_idle) def get_state(self): return self._state def get_state_string(self): return {Monitor.STATE_INACTIVE: "inactive", Monitor.STATE_IDLE: "idle", Monitor.STATE_PRE_ALARM: "pre-alarm", Monitor.STATE_ALARM: "alarm", Monitor.STATE_POST_ALARM: "post-alarm", }.get(self._state, "unknown") def _set_state(self, state): oldstate = self._state self._state = state if oldstate != state: self.emit("state-changed", oldstate, state) logging.debug("monitor '%s' state changed to '%s'" %(self.get_name(), self.get_state_string())) if state == Monitor.STATE_INACTIVE: logging.warning("monitor '%s' inactive"%self.get_name()) def get_pipeline(self): return self._pipeline def get_source(self): return self._source def set_source(self, source): self._source = source def get_alarm_pre_time(self): return self._alarm_pre_time def set_alarm_pre_time(self, alarm_pre_time): self._alarm_pre_time = alarm_pre_time def get_alarm_post_time(self): return self._alarm_post_time def set_alarm_post_time(self, alarm_post_time): self._alarm_post_time = alarm_post_time def get_geometry(self): return self._geometry def set_geometry(self, geometry): if geometry: geom = (int(geometry[0]), int(geometry[1])) if geom[0] <= 0 or geom[0] <= 0: raise ValueError("invalid geometry %d x %d"%(geom[0], geom[1])) self._geometry = geom else: self._geometry = None def get_framerate(self): return self._framerate def set_framerate(self, framerate): if not framerate: self._framerate = None else: if framerate < 1 or framerate > 255: raise ValueError("frame rate out of range, %s"%framerate) self._framerate = int(framerate) def get_debug(self): return self._debug def set_debug(self, debug): s = str(debug).lower().strip() if s == "yes" or s == "true" or (s.isdigit() and s != "0"): self._debug = True else: self._debug = False def get_launch_source_elements(self): if not self._source: raise ValueError("no source") element = get_element_from_url(self._source) if not element: raise ValueError("cannot handl URL '%s'"%source) elements = [element] if element.startswith("v4l2src"): if self._geometry or self._framerate: caps = "capsfilter caps=video/x-raw-yuv" if self._geometry: caps += ",width=%d,height=%d"%self._geometry if self._framerate: caps += ",framerate=%d/1"%self._framerate elements.append(caps) else: elements.append("decodebin") if self._geometry: elements.append("videoscale") elements.append("video/x-raw-yuv,width=%d,height=%d" %self._geometry) if self._framerate: elements.append("videorate") elements.append("video/x-raw-yuv,framerate=%d/1" %self._framerate) return elements def get_launch_motion_elements(self): elements = [] elements.append("ffmpegcolorspace") elements.append("motion debug=%s min-change=4" %(self._debug and "yes" or "no")) return elements def get_launch_destination_elements(self): elements = [] #elements.append("tee name=t") #elements.append("queue") #elements.append("ffmpegcolorspace") #elements.append("xvimagesink t.") elements.append("queue") elements.append("appsink name=appsink emit-signals=true t.") return elements def get_launch_string(self): elements = self.get_launch_source_elements() elements += self.get_launch_motion_elements() elements += self.get_launch_destination_elements() return " ! ".join(elements) def get_name(self): return self._name def restart(self): self.stop() self.start() def start(self): if not self._pipeline: self._pipeline = gst.parse_launch(self.get_launch_string()) appsink = self._pipeline.get_by_name("appsink") if appsink: appsink.connect("new-buffer", self.on_new_buffer) else: logging.warning("no element named 'appsink' found in pipeline, disabling recording") bus = self._pipeline.get_bus() bus.add_signal_watch() bus.connect("message", self.on_message) logging.debug("created new pipeline for monitor '%s'" %self.get_name()) if self._pipeline and self._pipeline.get_state() != gst.STATE_PLAYING: self._pipeline.set_state(gst.STATE_PLAYING) def stop(self): if self._pipeline and self_pipeline.get_state() != gst.STATE_NULL: self._pipeline.set_state(gst.STATE_NULL) self._pipeline = None def pause(self): if self._pipeline and self._pipeline.get_state() == gst.STATE_PLAYING: self._pipeline.set_state(gst.STATE_PAUSING) def on_message(self, bus, message): t = message.type if t == gst.MESSAGE_EOS: if self._pipeline: self._pipeline.set_state(gst.STATE_NULL) self._set_state(Monitor.STATE_INACTIVE) elif t == gst.MESSAGE_ERROR: if self._pipeline: self._pipeline.set_state(gst.STATE_NULL) self._set_state(Monitor.STATE_INACTIVE) err, debug = message.parse_error() error = str(err) if debug: error += " (%s)"%debug logging.warning("monitor '%s' received error; %s"%(self.get_name(), error)) elif t == gst.MESSAGE_STATE_CHANGED: if self._pipeline.get_state() == gst.STATE_NULL: self._set_state(Monitor.STATE_INACTIVE) else: self._set_state(Monitor.STATE_IDLE) elif t == gst.MESSAGE_ELEMENT: if message.structure.get_name() == "motion": if self._state != Monitor.STATE_INACTIVE: d = dict(message.structure) if d.get("start"): if self._state == Monitor.STATE_IDLE: self._alarm_timestamp = time.time() self._set_state(Monitor.STATE_PRE_ALARM) else: self._set_state(Monitor.STATE_ALARM) else: if self._state == Monitor.STATE_ALARM: self._alarm_timestamp = time.time() self._set_state(Monitor.STATE_POST_ALARM) else: self._set_state(Monitor.STATE_IDLE) def on_idle(self): if self._state == Monitor.STATE_PRE_ALARM: if time.time() - self._alarm_timestamp > self._alarm_pre_time: self._set_state(Monitor.STATE_ALARM) elif self._state == Monitor.STATE_POST_ALARM: if time.time() - self._alarm_timestamp > self._alarm_post_time: self._set_state(Monitor.STATE_IDLE) return True def on_new_buffer(self, appsink): buf = appsink.emit("pull-buffer") print len(buf.data) class Monitors(object): def __init__(self, **kwargs): self._monitors = [] self._anon_count = 0 self._open_files = {} self._loop = gobject.MainLoop() def get_main_loop(self): return self._loop def append_instance(self, monitor): mon_names = {} if not monitor.get_name(): self._anon_count += 1 monitor._name = "Monitor %d"%(self._anon_count) while monitor.get_name() in (m.get_name() for m in self._monitors): monitor._name += "_" monitor.connect("state-changed", self.on_monitor_changed) self._monitors.append(monitor) logging.debug("gst-launch %s"%monitor.get_launch_string()) monitor.start() logging.info("added monitor '%s' using source '%s'"% (monitor.get_name(), monitor.get_source())) def append(self, monitor_string): name, source, d = get_monitor_arguments_from_string(arg) self.append_instance(Monitor(name, source, **d)) def remove_instance(self, monitor): if monitor in self._monitors: self._monitor.stop() self._monitors.remove(monitor) def remove(self, name): mon = self.get_monitor_by_name(name) if not mon: raise ValueError("invalid monitor '%s'"%name) self.remove_instance(mon) def start(self): self._loop.run() def stop(self): self._loop.quit() def get_is_running(self): return self._loop.is_running() def get_monitor_by_name(self, name): for mon in self._monitors: if mon.get_name() == name: return mon return None def restart_monitor(self, name): mon = self.get_monitor_by_name(name) if not mon: raise ValueError("invalid monitor '%s'"%name) mon.restart() def on_monitor_changed(self, monitor, old_state, new_state): pass class DBusCommunicator(dbus.service.Object): def __init__(self, monitors): dbus.service.Object.__init__(self, dbus.SessionBus(), "/gstreamer/motion", ) self.monitors = monitors @dbus.service.method(dbus_interface='gstreamer.Motion', out_signature='s') def get_status(self): return "ok" if __name__ == "__main__": from optparse import OptionParser parser = OptionParser() parser.add_option("-D", "--debug", action="store_true", dest="debug", default=False, help="render debug to buffer and print messages") (options, args) = parser.parse_args() logging.basicConfig(format="%(asctime)s %(levelname)s: %(message)s", level=logging.DEBUG) monitors = Monitors() for arg in args: monitors.append(arg) communicator = DBusCommunicator(monitors) monitors.start()