Subversion Repositories HomeAutomation

Rev

Go to most recent revision | Blame | Compare with Previous | Last modification | View Log | SVN | RSS feed

  1. ############################################################################
  2. #    Copyright (C) 2008 by Mattias Runge   #
  3. #    mattias@runge.se   #
  4. #                                                                          #
  5. #    This program is free software; you can redistribute it and#or modify  #
  6. #    it under the terms of the GNU General Public License as published by  #
  7. #    the Free Software Foundation; either version 2 of the License, or     #
  8. #    (at your option) any later version.                                   #
  9. #                                                                          #
  10. #    This program is distributed in the hope that it will be useful,       #
  11. #    but WITHOUT ANY WARRANTY; without even the implied warranty of        #
  12. #    MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.  See the         #
  13. #    GNU General Public License for more details.                          #
  14. #                                                                          #
  15. #    You should have received a copy of the GNU General Public License     #
  16. #    along with this program; if not, write to the                         #
  17. #    Free Software Foundation, Inc.,                                       #
  18. #    59 Temple Place - Suite 330, Boston, MA  02111-1307, USA.             #
  19. ############################################################################
  20. import os
  21. import sys
  22. import asyncore
  23. import string
  24. import socket
  25. import Queue
  26. import threading
  27. import time
  28. import ServiceManager
  29. import CanModuleTypes
  30. import Util
  31. import Settings
  32.  
  33. gBackendCan = False
  34.  
  35. def GetBackendCan():
  36.     global gBackendCan
  37.     if gBackendCan == False:
  38.         gBackendCan = BackendCan(Settings.Settings["CANDAEMON_ADDRESS"], int(Settings.Settings["CANDAEMON_PORT"]))
  39.     return gBackendCan
  40.  
  41.  
  42.  
  43. def ConvertHardwareIdToString(hwid):
  44.     return str("%X" % hwid[0]).zfill(2) + str("%X" % hwid[1]).zfill(2) + str("%X" % hwid[2]).zfill(2) + str("%X" % hwid[3]).zfill(2)
  45.  
  46. class CanMessage:
  47.     def __init__(self):
  48.         self.ClassName = 0
  49.         self.Direction = 0
  50.         self.ModuleType = 0
  51.         self.ModuleId = 0
  52.         self.Command = 0
  53.         self.Data = []
  54.        
  55.     def Parse(self, msg):
  56.         parts = msg.split(' ')
  57.  
  58.         id_bits = Util.hex2bin(parts[1]);
  59.        
  60.         self.ClassName = int(id_bits[4:7], 2)
  61.         self.Direction = int(id_bits[7:8], 2)
  62.         self.ModuleType = int(id_bits[8:16], 2)
  63.         self.ModuleId = int(id_bits[16:24], 2)
  64.         self.Command = int(id_bits[24:32], 2)
  65.  
  66.         count = 4
  67.         while count < len(parts):
  68.             self.Data.append(int(Util.hex2bin(parts[count]), 2))
  69.             count += 1
  70.            
  71.     def Compile(self):
  72.         bClassName = Util.int2bin(int(self.ClassName)).zfill(4)
  73.         bDirection = Util.int2bin(int(self.Direction)).zfill(1)
  74.         bModuleType = Util.int2bin(int(self.ModuleType)).zfill(8)
  75.         bModuleId = Util.int2bin(int(self.ModuleId)).zfill(8)
  76.         bCommand = Util.int2bin(int(self.Command)).zfill(8)
  77.        
  78.         id_bits = "000" + bClassName + bDirection + bModuleType + bModuleId + bCommand
  79.        
  80.         id = Util.bin2hex(id_bits)
  81.        
  82.         bytes = ""
  83.         for byte in self.Data:
  84.             bytes += Util.bin2hex(Util.int2bin(byte).zfill(8)) + " "
  85.        
  86.         cmd = "PKT " + id + " 1 0 " + bytes
  87.         return cmd.strip(" ")
  88.  
  89. class CanNode ():
  90.     HardwareId = [0,0,0,0]
  91.     LastOnline = 0
  92.     Services = {}
  93.    
  94.     def __init__(self, hwid):
  95.         self.Services = {}
  96.         self.HardwareId = hwid
  97.         self.UpdateLastOnline()
  98.         self.canBackend = GetBackendCan()
  99.         self.serviceMan = ServiceManager.GetServiceManager()
  100.         self.AskForModules()
  101.    
  102.     def AskForModules(self):
  103.         canMsg = CanMessage()
  104.         canMsg.ClassName = 11 # 0x0B CAN_CLASS_MODULE_NMT
  105.         canMsg.ModuleType = 0
  106.         canMsg.ModuleId = 0
  107.         canMsg.Direction = 0
  108.         canMsg.Command = 0 # 0x00 CAN_CMD_MODULE_NMT_LIST
  109.        
  110.         canMsg.Data = self.HardwareId
  111.        
  112.         print "CAN: " + self.GetHardwareIdString() + ": Asking for modules..."
  113.        
  114.         self.canBackend.QueueMsg(canMsg.Compile())
  115.    
  116.     def AddModule(self, ModuleType, ModuleId):
  117.  
  118.         try:
  119.             serviceType = CanModuleTypes.Types[ModuleType]
  120.             serviceId = "CAN_" + str(ModuleType) + "_" + str(ModuleId)
  121.  
  122.             if not self.serviceMan.HasService(serviceId, serviceType):
  123.                 self.serviceMan.AddService(serviceId, serviceType)
  124.                
  125.                 print "CAN: " + self.GetHardwareIdString() + ": Added " + serviceType + " service with module id " + hex(ModuleId)
  126.             else:
  127.                 print "CAN: " + self.GetHardwareIdString() + ": Service " + serviceType + " already started, module id " + hex(ModuleId)
  128.  
  129.             service = self.serviceMan.GetService(serviceId, serviceType)
  130.            
  131.             service.ModuleType = ModuleType
  132.             service.ModuleId = ModuleId
  133.             service.SetOnline()
  134.  
  135.             self.Services[str(ModuleType) + "_" + str(ModuleId)] = service;
  136.         except (KeyboardInterrupt, SystemExit):
  137.             raise
  138.         except KeyError:
  139.             print "CAN: " + self.GetHardwareIdString() + ": Unknown module with id " + hex(ModuleId) + ", type " + hex(ModuleType)
  140.            
  141.    
  142.     def UpdateLastOnline(self):
  143.         self.LastOnline = time.time()
  144.         for key in self.Services:
  145.             self.Services[key].SetOnline()
  146.  
  147.     def CheckOnline(self):
  148.         if self.LastOnline+20 < time.time():
  149.             self.SetServicesOffline()
  150.  
  151.     def SetServicesOffline(self):
  152.         print "CAN: " + self.GetHardwareIdString() + ": Node is offline, stopping services"
  153.         for key in self.Services:
  154.             self.Services[key].SetOffline()
  155.  
  156.     def GetHardwareIdString(self):
  157.         return ConvertHardwareIdToString(self.HardwareId)
  158.        
  159.  
  160.  
  161. class BackendCan (asyncore.dispatcher, threading.Thread):
  162.    
  163.     nodes = {}
  164.     msgCallbacks = []
  165.    
  166.     def __init__(self, serverHost, serverPort):
  167.         self.nodes = {}
  168.         self.msgCallbacks = []
  169.  
  170.         threading.Thread.__init__(self)
  171.        
  172.         asyncore.dispatcher.__init__(self)
  173.        
  174.         self.serviceMan = ServiceManager.GetServiceManager()
  175.        
  176.         self.nodesLock = threading.Lock()
  177.  
  178.         self.queue = Queue.Queue()
  179.         self.queueLock = threading.Lock()
  180.         self.connected = False
  181.        
  182.         self.serverHost = serverHost;
  183.         self.serverPort = serverPort;
  184.  
  185.         self.timer = threading.Timer(20.0, self.CheckNodes)
  186.  
  187.         self.Connect()
  188.  
  189.         self.setDaemon(True)
  190.  
  191.         self.start()
  192.  
  193.     def Connect(self):
  194.         try:
  195.             print "CAN: Trying to connect to canDaemon at " + self.serverHost + ":" + str(self.serverPort)
  196.  
  197.             self.create_socket(socket.AF_INET, socket.SOCK_STREAM)
  198.  
  199.             self.connect((self.serverHost, self.serverPort))
  200.        
  201.         except (KeyboardInterrupt, SystemExit):
  202.             raise
  203.         except socket.gaierror, e:
  204.             print "CAN: Address-related error connecting to server: %s" % e
  205.             self.handle_close()
  206.         except socket.error, e:
  207.             print "CAN: Connection error: %s" % e
  208.             self.handle_close()
  209.  
  210.     def CheckNodes(self):
  211.         try:
  212.             print "CAN: Checking if nodes are offline"
  213.  
  214.             self.nodesLock.acquire()
  215.  
  216.             for key in self.nodes:
  217.                 self.nodes[key].CheckOnline()
  218.  
  219.             self.nodesLock.release()
  220.  
  221.             self.StartTimer()
  222.         except (KeyboardInterrupt, SystemExit):
  223.             raise
  224.  
  225.     def SetServicesOffline(self):
  226.         self.nodesLock.acquire()
  227.  
  228.         for key in self.nodes:
  229.             self.nodes[key].SetServicesOffline()
  230.  
  231.         self.nodesLock.release()
  232.  
  233.     def StartTimer(self):
  234.        
  235.         self.timer.start()
  236.  
  237.     def handle_connect(self):
  238.         pass
  239.  
  240.     def handle_close(self):
  241.         self.connected = False
  242.         print "CAN: " + "Connection to canDaemon closed"
  243.         self.close()
  244.         self.SetServicesOffline()
  245.         self.timer.cancel()
  246.  
  247.         print "CAN: Waiting for 10 seconds..."
  248.         timer = threading.Timer(10.0, self.Connect)
  249.         timer.start()
  250.         timer.join()
  251.        
  252.     def run(self):
  253.         #self.log('run')
  254.         asyncore.loop(1)
  255.        
  256.     def writable(self):
  257.         self.queueLock.acquire()
  258.         state = not self.queue.empty()
  259.         #self.log("writable: " + str(state))
  260.         self.queueLock.release()
  261.         return state
  262.    
  263.     def handle_write(self):
  264.         self.queueLock.acquire()
  265.        
  266.         while not self.queue.empty():
  267.             msg = self.queue.get()
  268.             self.send(msg + "\n")
  269.             #self.log("SENT: " + msg)
  270.             self.queue.task_done()
  271.             time.sleep(0.0001)
  272.            
  273.         self.queueLock.release()
  274.        
  275.     def handle_read(self):
  276.         try:
  277.             more = self.recv(1024)
  278.  
  279.             if not self.connected:
  280.                 self.connected = True
  281.                 print "CAN: " + "Connected to canDaemon"
  282.  
  283.                 self.StartTimer()
  284.  
  285.             if not more:
  286.                 self.handle_close()
  287.            
  288.             recvList = more.split("\n")
  289.        
  290.             for msg in recvList:
  291.                 self.handle_msg(msg)
  292.  
  293.         except (KeyboardInterrupt, SystemExit):
  294.             raise
  295.         except socket.error, e:
  296.             print "CAN: Error receiving data: %s" % e
  297.             self.handle_close()
  298.            
  299.     def handle_msg(self, msg):
  300.         if len(msg) > 0:
  301.             #self.log("RECV: " + msg)
  302.            
  303.             canMsg = CanMessage()
  304.             canMsg.Parse(msg)
  305.            
  306.             if canMsg.ClassName == 0: # 0x00 CAN_NMT
  307.                 if canMsg.ModuleType == 44: # 0x2C CAN_NMT_HEARTBEAT
  308.                     self.ProcessHeartbeat(canMsg.Data)
  309.             else:
  310.                 if canMsg.Command == 00: # 0x00 CAN_CMD_MODULE_NMT_LIST
  311.                     self.ProcessModuleList(canMsg)
  312.                 else:
  313.                     for callback in self.msgCallbacks:
  314.                         callback(canMsg)
  315.            
  316.     def ProcessModuleList(self, canMsg):
  317.         HardwareIdString = ConvertHardwareIdToString(canMsg.Data)
  318.        
  319.         self.nodesLock.acquire()
  320.  
  321.         if self.nodes.has_key(HardwareIdString):
  322.             #print "CAN: " + "Received module info from " + HardwareIdString
  323.             self.nodes[HardwareIdString].AddModule(canMsg.ModuleType, canMsg.ModuleId)
  324.         else:
  325.             print "CAN: " + "Received module info from " + HardwareIdString + " but that node is not online, ignoring"
  326.  
  327.         self.nodesLock.release()
  328.    
  329.     def ProcessHeartbeat(self, HardwareId):
  330.         HardwareIdString = ConvertHardwareIdToString(HardwareId)
  331.        
  332.         print "CAN: " + HardwareIdString + ": Received heartbeat"
  333.  
  334.         self.nodesLock.acquire()
  335.  
  336.         if self.nodes.has_key(HardwareIdString):
  337.             self.nodes[HardwareIdString].UpdateLastOnline()
  338.         else:
  339.             self.nodes[HardwareIdString] = CanNode(HardwareId)
  340.  
  341.         self.nodesLock.release()
  342.    
  343.     def QueueMsg(self, msg):
  344.         if self.connected:
  345.             self.queueLock.acquire()
  346.             self.queue.put(msg)
  347.             self.queueLock.release()
  348.             self.handle_write()
  349.         else:
  350.             self.queueLock.acquire()
  351.             self.queue.put(msg)
  352.             self.queueLock.release()
  353.        
  354.     def AddCallback(self, callback):
  355.         self.msgCallbacks.append(callback)
  356.