reorganize, build dockerfiles for ledserver and html-controller
This commit is contained in:
@@ -0,0 +1,729 @@
|
||||
'''
|
||||
The MIT License (MIT)
|
||||
Copyright (c) 2013 Dave P.
|
||||
https://github.com/dpallot/simple-websocket-server
|
||||
'''
|
||||
import sys
|
||||
VER = sys.version_info[0]
|
||||
import socketserver
|
||||
from http.server import BaseHTTPRequestHandler
|
||||
from io import StringIO, BytesIO
|
||||
|
||||
import hashlib
|
||||
import base64
|
||||
import socket
|
||||
import struct
|
||||
import ssl
|
||||
import errno
|
||||
import codecs
|
||||
from collections import deque
|
||||
from select import select
|
||||
|
||||
import traceback
|
||||
import time
|
||||
|
||||
__all__ = ['WebSocket',
|
||||
'SimpleWebSocketServer',
|
||||
'SimpleSSLWebSocketServer']
|
||||
|
||||
def _check_unicode(val):
|
||||
return isinstance(val, str)
|
||||
|
||||
class HTTPRequest(BaseHTTPRequestHandler):
|
||||
def __init__(self, request_text):
|
||||
if VER >= 3:
|
||||
self.rfile = BytesIO(request_text)
|
||||
else:
|
||||
self.rfile = StringIO(request_text)
|
||||
self.raw_requestline = self.rfile.readline()
|
||||
self.error_code = self.error_message = None
|
||||
self.parse_request()
|
||||
|
||||
_VALID_STATUS_CODES = [1000, 1001, 1002, 1003, 1007, 1008,
|
||||
1009, 1010, 1011, 3000, 3999, 4000, 4999]
|
||||
|
||||
HANDSHAKE_STR = (
|
||||
"HTTP/1.1 101 Switching Protocols\r\n"
|
||||
"Upgrade: WebSocket\r\n"
|
||||
"Connection: Upgrade\r\n"
|
||||
"Sec-WebSocket-Accept: %(acceptstr)s\r\n\r\n"
|
||||
)
|
||||
|
||||
GUID_STR = '258EAFA5-E914-47DA-95CA-C5AB0DC85B11'
|
||||
|
||||
STREAM = 0x0
|
||||
TEXT = 0x1
|
||||
BINARY = 0x2
|
||||
CLOSE = 0x8
|
||||
PING = 0x9
|
||||
PONG = 0xA
|
||||
|
||||
HEADERB1 = 1
|
||||
HEADERB2 = 3
|
||||
LENGTHSHORT = 4
|
||||
LENGTHLONG = 5
|
||||
MASK = 6
|
||||
PAYLOAD = 7
|
||||
|
||||
MAXHEADER = 65536
|
||||
MAXPAYLOAD = 33554432
|
||||
|
||||
class WebSocket(object):
|
||||
|
||||
def __init__(self, server, sock, address):
|
||||
self.server = server
|
||||
self.client = sock
|
||||
self.address = address
|
||||
|
||||
self.handshaked = False
|
||||
self.headerbuffer = bytearray()
|
||||
self.headertoread = 2048
|
||||
|
||||
self.fin = 0
|
||||
self.data = bytearray()
|
||||
self.opcode = 0
|
||||
self.hasmask = 0
|
||||
self.maskarray = None
|
||||
self.length = 0
|
||||
self.lengtharray = None
|
||||
self.index = 0
|
||||
self.request = None
|
||||
self.usingssl = False
|
||||
self.lastping = 0
|
||||
|
||||
self.frag_start = False
|
||||
self.frag_type = BINARY
|
||||
self.frag_buffer = None
|
||||
self.frag_decoder = codecs.getincrementaldecoder('utf-8')(errors='strict')
|
||||
self.closed = False
|
||||
self.sendq = deque()
|
||||
|
||||
self.state = HEADERB1
|
||||
|
||||
# restrict the size of header and payload for security reasons
|
||||
self.maxheader = MAXHEADER
|
||||
self.maxpayload = MAXPAYLOAD
|
||||
|
||||
def handleMessage(self):
|
||||
"""
|
||||
Called when websocket frame is received.
|
||||
To access the frame data call self.data.
|
||||
|
||||
If the frame is Text then self.data is a unicode object.
|
||||
If the frame is Binary then self.data is a bytearray object.
|
||||
"""
|
||||
pass
|
||||
|
||||
def handleConnected(self):
|
||||
"""
|
||||
Called when a websocket client connects to the server.
|
||||
"""
|
||||
pass
|
||||
|
||||
def handleClose(self):
|
||||
"""
|
||||
Called when a websocket server gets a Close frame from a client.
|
||||
"""
|
||||
pass
|
||||
|
||||
def _handlePacket(self):
|
||||
if self.opcode == CLOSE:
|
||||
pass
|
||||
elif self.opcode == STREAM:
|
||||
pass
|
||||
elif self.opcode == TEXT:
|
||||
pass
|
||||
elif self.opcode == BINARY:
|
||||
pass
|
||||
elif self.opcode == PONG or self.opcode == PING:
|
||||
self.lastping = time.time()
|
||||
if len(self.data) > 125:
|
||||
print('control frame length can not be > 125')
|
||||
raise Exception('control frame length can not be > 125')
|
||||
else:
|
||||
# unknown or reserved opcode so just close
|
||||
print('unknown opcode')
|
||||
raise Exception('unknown opcode')
|
||||
|
||||
if self.opcode == CLOSE:
|
||||
status = 1000
|
||||
reason = u''
|
||||
length = len(self.data)
|
||||
|
||||
if length == 0:
|
||||
pass
|
||||
elif length >= 2:
|
||||
status = struct.unpack_from('!H', self.data[:2])[0]
|
||||
reason = self.data[2:]
|
||||
|
||||
if status not in _VALID_STATUS_CODES:
|
||||
status = 1002
|
||||
|
||||
if len(reason) > 0:
|
||||
try:
|
||||
reason = reason.decode('utf8', errors='strict')
|
||||
except:
|
||||
status = 1002
|
||||
else:
|
||||
status = 1002
|
||||
|
||||
self.close(status, reason)
|
||||
return
|
||||
|
||||
elif self.fin == 0:
|
||||
if self.opcode != STREAM:
|
||||
if self.opcode == PING or self.opcode == PONG:
|
||||
print('control messages can not be fragmented')
|
||||
raise Exception('control messages can not be fragmented')
|
||||
|
||||
self.frag_type = self.opcode
|
||||
self.frag_start = True
|
||||
self.frag_decoder.reset()
|
||||
|
||||
if self.frag_type == TEXT:
|
||||
self.frag_buffer = []
|
||||
utf_str = self.frag_decoder.decode(self.data, final = False)
|
||||
if utf_str:
|
||||
self.frag_buffer.append(utf_str)
|
||||
else:
|
||||
self.frag_buffer = bytearray()
|
||||
self.frag_buffer.extend(self.data)
|
||||
|
||||
else:
|
||||
if self.frag_start is False:
|
||||
print('fragmentation protocol error')
|
||||
raise Exception('fragmentation protocol error')
|
||||
|
||||
if self.frag_type == TEXT:
|
||||
utf_str = self.frag_decoder.decode(self.data, final = False)
|
||||
if utf_str:
|
||||
self.frag_buffer.append(utf_str)
|
||||
else:
|
||||
self.frag_buffer.extend(self.data)
|
||||
|
||||
else:
|
||||
if self.opcode == STREAM:
|
||||
if self.frag_start is False:
|
||||
print('fragmentation protocol error')
|
||||
raise Exception('fragmentation protocol error')
|
||||
|
||||
if self.frag_type == TEXT:
|
||||
utf_str = self.frag_decoder.decode(self.data, final = True)
|
||||
self.frag_buffer.append(utf_str)
|
||||
self.data = u''.join(self.frag_buffer)
|
||||
else:
|
||||
self.frag_buffer.extend(self.data)
|
||||
self.data = self.frag_buffer
|
||||
|
||||
self.handleMessage()
|
||||
|
||||
self.frag_decoder.reset()
|
||||
self.frag_type = BINARY
|
||||
self.frag_start = False
|
||||
self.frag_buffer = None
|
||||
|
||||
elif self.opcode == PING:
|
||||
self._sendMessage(False, PONG, self.data)
|
||||
|
||||
elif self.opcode == PONG:
|
||||
pass
|
||||
|
||||
else:
|
||||
if self.frag_start is True:
|
||||
print('fragmentation protocol error')
|
||||
raise Exception('fragmentation protocol error')
|
||||
|
||||
if self.opcode == TEXT:
|
||||
try:
|
||||
self.data = self.data.decode('utf8', errors='strict')
|
||||
except Exception as exp:
|
||||
print('invalid utf-8 payload')
|
||||
raise Exception('invalid utf-8 payload')
|
||||
|
||||
self.handleMessage()
|
||||
|
||||
|
||||
def _handleData(self):
|
||||
# do the HTTP header and handshake
|
||||
if self.handshaked is False:
|
||||
|
||||
data = self.client.recv(self.headertoread)
|
||||
if not data:
|
||||
print('remote socket closed')
|
||||
raise Exception('remote socket closed')
|
||||
|
||||
else:
|
||||
# accumulate
|
||||
self.headerbuffer.extend(data)
|
||||
|
||||
if len(self.headerbuffer) >= self.maxheader:
|
||||
print('header exceeded allowable size')
|
||||
raise Exception('header exceeded allowable size')
|
||||
|
||||
# indicates end of HTTP header
|
||||
if b'\r\n\r\n' in self.headerbuffer:
|
||||
self.request = HTTPRequest(self.headerbuffer)
|
||||
|
||||
# handshake rfc 6455
|
||||
try:
|
||||
key = self.request.headers['Sec-WebSocket-Key']
|
||||
k = key.encode('ascii') + GUID_STR.encode('ascii')
|
||||
k_s = base64.b64encode(hashlib.sha1(k).digest()).decode('ascii')
|
||||
hStr = HANDSHAKE_STR % {'acceptstr': k_s}
|
||||
self.sendq.append((BINARY, hStr.encode('ascii')))
|
||||
self.handshaked = True
|
||||
self.handleConnected()
|
||||
except Exception as e:
|
||||
print(e,traceback.format_exc())
|
||||
print('handshake failed: %s', str(e))
|
||||
raise Exception('handshake failed: %s', str(e))
|
||||
|
||||
# else do normal data
|
||||
else:
|
||||
data = self.client.recv(16384)
|
||||
if not data:
|
||||
print("remote socket closed")
|
||||
raise Exception("remote socket closed")
|
||||
|
||||
if VER >= 3:
|
||||
for d in data:
|
||||
self._parseMessage(d)
|
||||
else:
|
||||
for d in data:
|
||||
self._parseMessage(ord(d))
|
||||
|
||||
def close(self, status = 1000, reason = u''):
|
||||
"""
|
||||
Send Close frame to the client. The underlying socket is only closed
|
||||
when the client acknowledges the Close frame.
|
||||
|
||||
status is the closing identifier.
|
||||
reason is the reason for the close.
|
||||
"""
|
||||
try:
|
||||
if self.closed is False:
|
||||
close_msg = bytearray()
|
||||
close_msg.extend(struct.pack("!H", status))
|
||||
if _check_unicode(reason):
|
||||
close_msg.extend(reason.encode('utf-8'))
|
||||
else:
|
||||
close_msg.extend(reason)
|
||||
|
||||
self._sendMessage(False, CLOSE, close_msg)
|
||||
|
||||
finally:
|
||||
self.closed = True
|
||||
|
||||
|
||||
def _sendBuffer(self, buff, send_all = False):
|
||||
size = len(buff)
|
||||
tosend = size
|
||||
already_sent = 0
|
||||
|
||||
while tosend > 0:
|
||||
try:
|
||||
# i should be able to send a bytearray
|
||||
sent = self.client.send(buff[already_sent:])
|
||||
if sent == 0:
|
||||
raise RuntimeError('socket connection broken')
|
||||
|
||||
already_sent += sent
|
||||
tosend -= sent
|
||||
|
||||
except socket.error as e:
|
||||
print(e,traceback.format_exc())
|
||||
# if we have full buffers then wait for them to drain and try again
|
||||
if e.errno in [errno.EAGAIN, errno.EWOULDBLOCK]:
|
||||
if send_all:
|
||||
continue
|
||||
return buff[already_sent:]
|
||||
else:
|
||||
print(e,traceback.format_exc())
|
||||
raise e
|
||||
|
||||
return None
|
||||
|
||||
def sendFragmentStart(self, data):
|
||||
"""
|
||||
Send the start of a data fragment stream to a websocket client.
|
||||
Subsequent data should be sent using sendFragment().
|
||||
A fragment stream is completed when sendFragmentEnd() is called.
|
||||
|
||||
If data is a unicode object then the frame is sent as Text.
|
||||
If the data is a bytearray object then the frame is sent as Binary.
|
||||
"""
|
||||
opcode = BINARY
|
||||
if _check_unicode(data):
|
||||
opcode = TEXT
|
||||
self._sendMessage(True, opcode, data)
|
||||
|
||||
def sendFragment(self, data):
|
||||
"""
|
||||
see sendFragmentStart()
|
||||
|
||||
If data is a unicode object then the frame is sent as Text.
|
||||
If the data is a bytearray object then the frame is sent as Binary.
|
||||
"""
|
||||
self._sendMessage(True, STREAM, data)
|
||||
|
||||
def sendFragmentEnd(self, data):
|
||||
"""
|
||||
see sendFragmentEnd()
|
||||
|
||||
If data is a unicode object then the frame is sent as Text.
|
||||
If the data is a bytearray object then the frame is sent as Binary.
|
||||
"""
|
||||
self._sendMessage(False, STREAM, data)
|
||||
|
||||
def sendMessage(self, data):
|
||||
"""
|
||||
Send websocket data frame to the client.
|
||||
|
||||
If data is a unicode object then the frame is sent as Text.
|
||||
If the data is a bytearray object then the frame is sent as Binary.
|
||||
"""
|
||||
opcode = BINARY
|
||||
if _check_unicode(data):
|
||||
opcode = TEXT
|
||||
self._sendMessage(False, opcode, data)
|
||||
|
||||
|
||||
def _sendMessage(self, fin, opcode, data):
|
||||
|
||||
payload = bytearray()
|
||||
|
||||
b1 = 0
|
||||
b2 = 0
|
||||
if fin is False:
|
||||
b1 |= 0x80
|
||||
b1 |= opcode
|
||||
|
||||
if _check_unicode(data):
|
||||
data = data.encode('utf-8')
|
||||
|
||||
length = len(data)
|
||||
payload.append(b1)
|
||||
|
||||
if length <= 125:
|
||||
b2 |= length
|
||||
payload.append(b2)
|
||||
|
||||
elif length >= 126 and length <= 65535:
|
||||
b2 |= 126
|
||||
payload.append(b2)
|
||||
payload.extend(struct.pack("!H", length))
|
||||
|
||||
else:
|
||||
b2 |= 127
|
||||
payload.append(b2)
|
||||
payload.extend(struct.pack("!Q", length))
|
||||
|
||||
if length > 0:
|
||||
payload.extend(data)
|
||||
|
||||
self.sendq.append((opcode, payload))
|
||||
|
||||
|
||||
def _parseMessage(self, byte):
|
||||
# read in the header
|
||||
if self.state == HEADERB1:
|
||||
|
||||
self.fin = byte & 0x80
|
||||
self.opcode = byte & 0x0F
|
||||
self.state = HEADERB2
|
||||
|
||||
self.index = 0
|
||||
self.length = 0
|
||||
self.lengtharray = bytearray()
|
||||
self.data = bytearray()
|
||||
|
||||
rsv = byte & 0x70
|
||||
if rsv != 0:
|
||||
print('RSV bit must be 0')
|
||||
raise Exception('RSV bit must be 0')
|
||||
|
||||
elif self.state == HEADERB2:
|
||||
mask = byte & 0x80
|
||||
length = byte & 0x7F
|
||||
|
||||
if self.opcode == PING and length > 125:
|
||||
print('ping packet is too large')
|
||||
raise Exception('ping packet is too large')
|
||||
|
||||
if mask == 128:
|
||||
self.hasmask = True
|
||||
else:
|
||||
self.hasmask = False
|
||||
|
||||
if length <= 125:
|
||||
self.length = length
|
||||
|
||||
# if we have a mask we must read it
|
||||
if self.hasmask is True:
|
||||
self.maskarray = bytearray()
|
||||
self.state = MASK
|
||||
else:
|
||||
# if there is no mask and no payload we are done
|
||||
if self.length <= 0:
|
||||
try:
|
||||
self._handlePacket()
|
||||
finally:
|
||||
self.state = HEADERB1
|
||||
self.data = bytearray()
|
||||
|
||||
# we have no mask and some payload
|
||||
else:
|
||||
#self.index = 0
|
||||
self.data = bytearray()
|
||||
self.state = PAYLOAD
|
||||
|
||||
elif length == 126:
|
||||
self.lengtharray = bytearray()
|
||||
self.state = LENGTHSHORT
|
||||
|
||||
elif length == 127:
|
||||
self.lengtharray = bytearray()
|
||||
self.state = LENGTHLONG
|
||||
|
||||
|
||||
elif self.state == LENGTHSHORT:
|
||||
self.lengtharray.append(byte)
|
||||
|
||||
if len(self.lengtharray) > 2:
|
||||
print('short length exceeded allowable size')
|
||||
raise Exception('short length exceeded allowable size')
|
||||
|
||||
if len(self.lengtharray) == 2:
|
||||
self.length = struct.unpack_from('!H', self.lengtharray)[0]
|
||||
|
||||
if self.hasmask is True:
|
||||
self.maskarray = bytearray()
|
||||
self.state = MASK
|
||||
else:
|
||||
# if there is no mask and no payload we are done
|
||||
if self.length <= 0:
|
||||
try:
|
||||
self._handlePacket()
|
||||
finally:
|
||||
self.state = HEADERB1
|
||||
self.data = bytearray()
|
||||
|
||||
# we have no mask and some payload
|
||||
else:
|
||||
#self.index = 0
|
||||
self.data = bytearray()
|
||||
self.state = PAYLOAD
|
||||
|
||||
elif self.state == LENGTHLONG:
|
||||
|
||||
self.lengtharray.append(byte)
|
||||
|
||||
if len(self.lengtharray) > 8:
|
||||
print('long length exceeded allowable size')
|
||||
raise Exception('long length exceeded allowable size')
|
||||
|
||||
if len(self.lengtharray) == 8:
|
||||
self.length = struct.unpack_from('!Q', self.lengtharray)[0]
|
||||
|
||||
if self.hasmask is True:
|
||||
self.maskarray = bytearray()
|
||||
self.state = MASK
|
||||
else:
|
||||
# if there is no mask and no payload we are done
|
||||
if self.length <= 0:
|
||||
try:
|
||||
self._handlePacket()
|
||||
finally:
|
||||
self.state = HEADERB1
|
||||
self.data = bytearray()
|
||||
|
||||
# we have no mask and some payload
|
||||
else:
|
||||
#self.index = 0
|
||||
self.data = bytearray()
|
||||
self.state = PAYLOAD
|
||||
|
||||
# MASK STATE
|
||||
elif self.state == MASK:
|
||||
self.maskarray.append(byte)
|
||||
|
||||
if len(self.maskarray) > 4:
|
||||
print('mask exceeded allowable size')
|
||||
raise Exception('mask exceeded allowable size')
|
||||
|
||||
if len(self.maskarray) == 4:
|
||||
# if there is no mask and no payload we are done
|
||||
if self.length <= 0:
|
||||
try:
|
||||
self._handlePacket()
|
||||
finally:
|
||||
self.state = HEADERB1
|
||||
self.data = bytearray()
|
||||
|
||||
# we have no mask and some payload
|
||||
else:
|
||||
#self.index = 0
|
||||
self.data = bytearray()
|
||||
self.state = PAYLOAD
|
||||
|
||||
# PAYLOAD STATE
|
||||
elif self.state == PAYLOAD:
|
||||
if self.hasmask is True:
|
||||
self.data.append( byte ^ self.maskarray[self.index % 4] )
|
||||
else:
|
||||
self.data.append( byte )
|
||||
|
||||
# if length exceeds allowable size then we except and remove the connection
|
||||
if len(self.data) >= self.maxpayload:
|
||||
print('payload exceeded allowable size')
|
||||
raise Exception('payload exceeded allowable size')
|
||||
|
||||
# check if we have processed length bytes; if so we are done
|
||||
if (self.index+1) == self.length:
|
||||
try:
|
||||
self._handlePacket()
|
||||
finally:
|
||||
#self.index = 0
|
||||
self.state = HEADERB1
|
||||
self.data = bytearray()
|
||||
else:
|
||||
self.index += 1
|
||||
|
||||
|
||||
class SimpleWebSocketServer(object):
|
||||
def __init__(self, host, port, websocketclass, selectInterval = 0.1):
|
||||
self.websocketclass = websocketclass
|
||||
self.serversocket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||||
self.serversocket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
||||
self.serversocket.settimeout(5)
|
||||
self.serversocket.bind((host, port))
|
||||
self.serversocket.listen(5)
|
||||
self.selectInterval = selectInterval
|
||||
self.connections = {}
|
||||
self.listeners = [self.serversocket]
|
||||
|
||||
def _decorateSocket(self, sock):
|
||||
return sock
|
||||
|
||||
def _constructWebSocket(self, sock, address):
|
||||
return self.websocketclass(self, sock, address)
|
||||
|
||||
def close(self):
|
||||
self.serversocket.close()
|
||||
|
||||
for desc, conn in self.connections.items():
|
||||
conn.close()
|
||||
self._handleClose(conn)
|
||||
|
||||
def _handleClose(self, client):
|
||||
client.client.close()
|
||||
# only call handleClose when we have a successful websocket connection
|
||||
if client.handshaked:
|
||||
#try:
|
||||
client.handleClose()
|
||||
#except:
|
||||
# print("timeout?")
|
||||
# pass
|
||||
|
||||
def serveonce(self):
|
||||
writers = []
|
||||
for fileno in self.listeners:
|
||||
if fileno == self.serversocket:
|
||||
continue
|
||||
client = self.connections[fileno]
|
||||
if client.sendq:
|
||||
writers.append(fileno)
|
||||
|
||||
if self.selectInterval:
|
||||
rList, wList, xList = select(self.listeners, writers, self.listeners, self.selectInterval)
|
||||
else:
|
||||
rList, wList, xList = select(self.listeners, writers, self.listeners)
|
||||
|
||||
for ready in wList:
|
||||
client = self.connections[ready]
|
||||
try:
|
||||
while client.sendq:
|
||||
opcode, payload = client.sendq.popleft()
|
||||
remaining = client._sendBuffer(payload)
|
||||
if remaining is not None:
|
||||
client.sendq.appendleft((opcode, remaining))
|
||||
break
|
||||
else:
|
||||
if opcode == CLOSE:
|
||||
print('received client close')
|
||||
raise Exception('received client close')
|
||||
|
||||
except Exception as n:
|
||||
print(n,traceback.format_exc())
|
||||
self._handleClose(client)
|
||||
del self.connections[ready]
|
||||
self.listeners.remove(ready)
|
||||
|
||||
for ready in rList:
|
||||
if ready == self.serversocket:
|
||||
sock = None
|
||||
try:
|
||||
sock, address = self.serversocket.accept()
|
||||
newsock = self._decorateSocket(sock)
|
||||
newsock.setblocking(0)
|
||||
fileno = newsock.fileno()
|
||||
self.connections[fileno] = self._constructWebSocket(newsock, address)
|
||||
self.listeners.append(fileno)
|
||||
except Exception as n:
|
||||
print(n,traceback.format_exc())
|
||||
if sock is not None:
|
||||
sock.close()
|
||||
else:
|
||||
if ready not in self.connections:
|
||||
continue
|
||||
client = self.connections[ready]
|
||||
try:
|
||||
client._handleData()
|
||||
except Exception as n:
|
||||
print(n,traceback.format_exc())
|
||||
self._handleClose(client)
|
||||
del self.connections[ready]
|
||||
self.listeners.remove(ready)
|
||||
|
||||
for failed in xList:
|
||||
if failed == self.serversocket:
|
||||
self.close()
|
||||
print('server socket failed')
|
||||
raise Exception('server socket failed')
|
||||
else:
|
||||
if failed not in self.connections:
|
||||
continue
|
||||
client = self.connections[failed]
|
||||
self._handleClose(client)
|
||||
del self.connections[failed]
|
||||
self.listeners.remove(failed)
|
||||
|
||||
def serveforever(self):
|
||||
while True:
|
||||
self.serveonce()
|
||||
|
||||
class SimpleSSLWebSocketServer(SimpleWebSocketServer):
|
||||
|
||||
def __init__(self, host, port, websocketclass, certfile,
|
||||
keyfile, version = ssl.PROTOCOL_TLSv1, selectInterval = 0.1):
|
||||
|
||||
SimpleWebSocketServer.__init__(self, host, port,
|
||||
websocketclass, selectInterval)
|
||||
|
||||
self.context = ssl.SSLContext(version)
|
||||
self.context.load_cert_chain(certfile, keyfile)
|
||||
|
||||
def close(self):
|
||||
super(SimpleSSLWebSocketServer, self).close()
|
||||
|
||||
def _decorateSocket(self, sock):
|
||||
sslsock = self.context.wrap_socket(sock, server_side=True)
|
||||
return sslsock
|
||||
|
||||
def _constructWebSocket(self, sock, address):
|
||||
ws = self.websocketclass(self, sock, address)
|
||||
ws.usingssl = True
|
||||
return ws
|
||||
|
||||
def serveforever(self):
|
||||
super(SimpleSSLWebSocketServer, self).serveforever()
|
||||
@@ -0,0 +1,131 @@
|
||||
from BackendProvider.Helper.SimpleWebSocketServer import SimpleWebSocketServer, WebSocket
|
||||
from threading import Thread
|
||||
from functools import partial
|
||||
|
||||
class ThreadedWebSocketServer(Thread):
|
||||
def __init__(self,effectController,rgbStripController):
|
||||
Thread.__init__(self)
|
||||
self.effectController = effectController
|
||||
self.rgbStripController = rgbStripController
|
||||
self.daemon = True
|
||||
self.stopped = False
|
||||
self.start()
|
||||
|
||||
def run(self):
|
||||
server = SimpleWebSocketServer('', 8001, partial(HTTPWebSocketsHandler,self.effectController, self.rgbStripController))
|
||||
while not self.stopped:
|
||||
server.serveonce()
|
||||
print("ThreadedWebSocketServer stopped")
|
||||
|
||||
def stop(self):
|
||||
self.stopped = True
|
||||
|
||||
import json
|
||||
|
||||
from rgbUtils import effectControllerJsonHelper
|
||||
from rgbUtils import rgbStripControllerJsonHelper
|
||||
|
||||
import traceback
|
||||
import logging
|
||||
|
||||
|
||||
CLIENT_TYPE_CONTROLLER = 0
|
||||
CLIENT_TYPE_STRIPE = 1
|
||||
CLIENT_TYPE_RECORDER = 2
|
||||
|
||||
class HTTPWebSocketsHandler(WebSocket):
|
||||
|
||||
def __init__(self, effectController, rgbStripController, *args, **kwargs):
|
||||
self.effectController = effectController
|
||||
self.rgbStripController = rgbStripController
|
||||
self.client_type = CLIENT_TYPE_CONTROLLER
|
||||
self.rgbStrip = None
|
||||
super().__init__(*args, **kwargs)
|
||||
|
||||
def handleMessage(self):
|
||||
try:
|
||||
print(self.address, self.data)
|
||||
data = json.loads(self.data)
|
||||
# Client Registration on the Websocket Server
|
||||
# maybe it would be better to use a websocket server thread for each client type,
|
||||
# can be done in future if there is too much latency
|
||||
if "register_client_type" in data:
|
||||
# the controler type, add handler on RGBStripContoller and send the current state of the controller
|
||||
if int(data['register_client_type']) is CLIENT_TYPE_CONTROLLER:
|
||||
self.client_type = CLIENT_TYPE_CONTROLLER
|
||||
# add effectController onControllerChangeHandler to get changes in the effectController eg start/stop effects, parameter updates, moved strips
|
||||
# register rgbStripController onRGBStripRegistered/UnRegistered handler to get noticed about new rgbStrips is not necessary
|
||||
# since we will get noticed from the effectController when it added the rgbStrip to the offEffect
|
||||
self.effectController.addOnControllerChangeHandler(self.onChange)
|
||||
# register new Stripes
|
||||
elif int(data['register_client_type']) is CLIENT_TYPE_STRIPE and "client_name" in data:
|
||||
self.client_type = CLIENT_TYPE_STRIPE
|
||||
# registers the strip with websocket object and name. the onRGBStripValueUpdate(rgbStrip) is called by
|
||||
# by the rgbStrip when an effectThread updates it
|
||||
# the self.rgbStrip variable is used to unregister the strip only
|
||||
self.rgbStrip = self.rgbStripController.registerRGBStrip(data["client_name"],self.onRGBStripValueUpdate)
|
||||
# register new Audio Recorders
|
||||
elif int(data['register_client_type']) is CLIENT_TYPE_RECORDER:
|
||||
self.client_type = CLIENT_TYPE_RECORDER
|
||||
|
||||
# controller responses are handled by the effectControllerJsonHelper
|
||||
if self.client_type is CLIENT_TYPE_CONTROLLER:
|
||||
response = effectControllerJsonHelper.responseHandler(self.effectController, self.rgbStripController, data)
|
||||
self.sendMessage(
|
||||
json.dumps({
|
||||
'response': response
|
||||
})
|
||||
)
|
||||
return
|
||||
# the stripe should usualy not send any data, i do not know why it should...
|
||||
elif self.client_type is CLIENT_TYPE_STRIPE:
|
||||
return
|
||||
# audio recorder responses are handled by the effectControllerJsonHandler
|
||||
elif self.client_type is CLIENT_TYPE_RECORDER:
|
||||
return
|
||||
except Exception as e:
|
||||
print(e, traceback.format_exc())
|
||||
|
||||
# notice about connects in terminal, the client has to register itself, see handleMessage
|
||||
def handleConnected(self):
|
||||
print(self.address, 'connected')
|
||||
|
||||
# unregister the onChangeHandler
|
||||
# for now this function is not called when a client times out,
|
||||
# so they don't get unregistered. i mean there is no function that
|
||||
# is called when a client times out.
|
||||
def handleClose(self):
|
||||
if self.client_type is CLIENT_TYPE_CONTROLLER:
|
||||
self.effectController.removeOnControllerChangeHandler(self.onChange)
|
||||
elif self.client_type is CLIENT_TYPE_STRIPE:
|
||||
self.rgbStripController.unregisterRGBStrip(self.rgbStrip)
|
||||
elif self.client_type is CLIENT_TYPE_RECORDER:
|
||||
pass
|
||||
print(self.address, 'closed')
|
||||
|
||||
# called when there are changes that should be pushed to the client.
|
||||
# - the effectController: start / stop effects, move strip to effect, changing effect params
|
||||
# - the rgbStripController: add/removing strips
|
||||
# -> CLIENT_TYPE_CONTROLLER
|
||||
def onChange(self):
|
||||
if self.client_type is CLIENT_TYPE_CONTROLLER:
|
||||
self.sendMessage(
|
||||
json.dumps({
|
||||
'effects': effectControllerJsonHelper.getEffects(self.effectController),
|
||||
'rgbStrips': rgbStripControllerJsonHelper.getRGBStrips(self.rgbStripController),
|
||||
'effectThreads': effectControllerJsonHelper.getEffectThreads(self.effectController)
|
||||
})
|
||||
)
|
||||
return
|
||||
elif self.client_type is CLIENT_TYPE_STRIPE:
|
||||
return
|
||||
elif self.client_type is CLIENT_TYPE_RECORDER:
|
||||
return
|
||||
|
||||
# when a rgbStrip value is changed, send json data to client
|
||||
def onRGBStripValueUpdate(self,rgbStrip):
|
||||
self.sendMessage(
|
||||
json.dumps({
|
||||
'data': rgbStripControllerJsonHelper.getRGBData(rgbStrip)
|
||||
})
|
||||
)
|
||||
@@ -0,0 +1,137 @@
|
||||
import threading
|
||||
import socketserver
|
||||
import socket
|
||||
import traceback
|
||||
from time import sleep, time
|
||||
import json
|
||||
import struct
|
||||
|
||||
|
||||
class ThreadedUDPServer(threading.Thread):
|
||||
def __init__(self, effectController, rgbStripController):
|
||||
threading.Thread.__init__(self)
|
||||
self.effectController = effectController
|
||||
self.rgbStripController = rgbStripController
|
||||
self.daemon = True
|
||||
self.stopped = False
|
||||
self.start()
|
||||
self.udpClientGuardian = self.UDPClientGuardian()
|
||||
self.udpClientGuardian.start()
|
||||
UDPClients
|
||||
|
||||
def run(self):
|
||||
self.server = socketserver.UDPServer(('', 8002), UDPStripHandler)
|
||||
self.server.effectController = self.effectController
|
||||
self.server.rgbStripController = self.rgbStripController
|
||||
self.server.serve_forever()
|
||||
print("ThreadedUDPServer stopped")
|
||||
|
||||
def stop(self):
|
||||
self.udpClientGuardian.stop()
|
||||
self.server.shutdown()
|
||||
|
||||
# check last pings from clients, responds with pong and remove clients
|
||||
# when there is no answer after 2 seconds
|
||||
class UDPClientGuardian(threading.Thread):
|
||||
def __init__(self):
|
||||
threading.Thread.__init__(self)
|
||||
self.stopped = False
|
||||
|
||||
def run(self):
|
||||
while not self.stopped:
|
||||
for key in list(UDPClients.keys()):
|
||||
if UDPClients[key].lastping + 2 < time():
|
||||
UDPClients[key].handleClose()
|
||||
del UDPClients[key]
|
||||
sleep(0.5)
|
||||
|
||||
def stop(self):
|
||||
self.stopped = True
|
||||
|
||||
|
||||
CLIENT_TYPE_CONTROLLER = 0
|
||||
CLIENT_TYPE_STRIPE = 1
|
||||
CLIENT_TYPE_RECORDER = 2
|
||||
|
||||
UDPClients = {}
|
||||
|
||||
|
||||
class UDPStripHandler(socketserver.BaseRequestHandler):
|
||||
|
||||
def handle(self):
|
||||
# print(self.client_address)
|
||||
if self.client_address not in UDPClients:
|
||||
UDPClients[self.client_address] = UDPClient(
|
||||
self.client_address, self.server.effectController, self.server.rgbStripController)
|
||||
UDPClients[self.client_address].handle(self.request)
|
||||
|
||||
|
||||
class UDPClient():
|
||||
def __init__(self, client_address, effectController, rgbStripController):
|
||||
self.client_type = None
|
||||
self.rgbStrip = None
|
||||
self.socket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
|
||||
self.client_address = client_address
|
||||
self.effectController = effectController
|
||||
self.rgbStripController = rgbStripController
|
||||
self.sendToClientLock = False
|
||||
self.lastping = time()
|
||||
|
||||
def handle(self, request):
|
||||
|
||||
clientdata = request[0].decode()
|
||||
self.socket = request[1]
|
||||
#print(time(),"{} wrote:".format(self.client_address))
|
||||
#print(time(),"clientdata -> ", clientdata)
|
||||
# socket.sendto(bytes("pong","utf-8"),self.client_address)
|
||||
|
||||
try:
|
||||
data = clientdata.split(':')
|
||||
# print(data)
|
||||
# r:1:srg strip name
|
||||
if data[0] == "r" and int(data[1]) == CLIENT_TYPE_STRIPE and data[2] != None:
|
||||
self.client_type = CLIENT_TYPE_STRIPE
|
||||
# registers the strip with websocket object and name. the onRGBStripValueUpdate(rgbStrip) is called by
|
||||
# by the rgbStrip when an effectThread updates it
|
||||
# the self.rgbStrip variable is used to unregister the strip only
|
||||
self.rgbStrip = self.rgbStripController.registerRGBStrip(
|
||||
data[2], self.onRGBStripValueUpdate)
|
||||
# s:ping
|
||||
if data[0] == "s" and data[1] == "ping":
|
||||
# if we got a ping and the client has no client type defined, send status unregistered, so the client knows that he has to register
|
||||
if self.client_type is None and self.socket is not None:
|
||||
self.sendToClient('s:unregistered')
|
||||
self.lastping = time()
|
||||
if data[0] == "u" and self.client_type == CLIENT_TYPE_STRIPE:
|
||||
led = int(data[1])
|
||||
self.sendToClient('d:'+str(led)+':'+str(self.rgbStrip.red[led])+':'+str(
|
||||
self.rgbStrip.green[led])+':'+str(self.rgbStrip.blue[led])+'')
|
||||
except Exception as e:
|
||||
print(e, traceback.format_exc())
|
||||
|
||||
# unregister the onChangeHandler
|
||||
# for now this function is not called when a client times out,
|
||||
# so they don't get unregistered. i mean there is no function that
|
||||
# is called when a client times out.
|
||||
|
||||
def handleClose(self):
|
||||
if self.client_type is CLIENT_TYPE_STRIPE:
|
||||
self.rgbStripController.unregisterRGBStrip(self.rgbStrip)
|
||||
#print(self.client_address, 'closed')
|
||||
|
||||
# when a rgbStrip value is changed, send not json data but a formated string to client
|
||||
# d:[id off the LED, always 0 on RGB strips]:[red value 0-255]:[green value 0-255]:[blue value 0-255]
|
||||
def onRGBStripValueUpdate(self, rgbStrip, led=0):
|
||||
return # we send the update as requested, to prevent flooding the module
|
||||
#self.sendToClient('d:'+str(led)+':'+str(rgbStrip.red[led])+':'+str(
|
||||
# rgbStrip.green[led])+':'+str(rgbStrip.blue[led])+'')
|
||||
|
||||
def sendToClient(self, message):
|
||||
while self.sendToClientLock is True:
|
||||
sleep(1)
|
||||
self.sendToClientLock = True
|
||||
if self.socket is not None:
|
||||
self.socket.sendto(
|
||||
message.encode(), self.client_address
|
||||
)
|
||||
self.sendToClientLock = False
|
||||
@@ -0,0 +1,11 @@
|
||||
FROM python:3.7.2-alpine3.9
|
||||
|
||||
RUN echo "@community http://dl-cdn.alpinelinux.org/alpine/edge/community" >> /etc/apk/repositories
|
||||
RUN apk add --update --no-cache ca-certificates gcc g++ curl openblas-dev@community
|
||||
RUN ln -s /usr/include/locale.h /usr/include/xlocale.h
|
||||
|
||||
WORKDIR /usr/src/LEDServer
|
||||
RUN pip3 install --no-cache-dir numpy scipy
|
||||
|
||||
COPY . /usr/src/LEDServer
|
||||
CMD python3.7 -u /usr/src/LEDServer/LEDServer.py
|
||||
Executable
+65
@@ -0,0 +1,65 @@
|
||||
#!/usr/bin/python
|
||||
import os
|
||||
import sys
|
||||
import time
|
||||
import traceback
|
||||
|
||||
|
||||
def main():
|
||||
try:
|
||||
|
||||
running = True
|
||||
if(os.path.dirname(sys.argv[0]) is not ""):
|
||||
os.chdir(os.path.dirname(sys.argv[0]))
|
||||
|
||||
# i want some external providers the effects can interact with. for example the music reaction.
|
||||
# Idea is: eg a Pi with a soundcard processing input via pyaudio and sending this data to the server. an musicEffect is bind to this input and processing it.
|
||||
# to be as flexible as possible the client registers with a name(sting), a type(string) and the data as an dict. the effect filters these clients by type. jea?
|
||||
|
||||
# rgbStrips register themselves at the rgbStripContoller
|
||||
# the rgbStripController calls the backend Provider's onChange function
|
||||
# when there are new values for the strip
|
||||
from rgbUtils.rgbStripController import rgbStripController
|
||||
rgbStripController = rgbStripController()
|
||||
rgbStripController.start()
|
||||
|
||||
# the effectController handles the effects and pushes the values to the rgbStripContoller
|
||||
# it also calls the backendProvider's onChange function when there are changes made on the effects
|
||||
from rgbUtils.effectController import effectController
|
||||
effectController = effectController(rgbStripController)
|
||||
|
||||
# register effectControllers onRGBStripRegistered and onRGBStripUnregistered handler on the rgbStripContoller to detect added or removed strips
|
||||
rgbStripController.addOnRGBStripRegisteredHandler(
|
||||
effectController.onRGBStripRegistered)
|
||||
rgbStripController.addOnRGBStripUnRegisteredHandler(
|
||||
effectController.onRGBStripUnRegistered)
|
||||
|
||||
# this is a "Backend Provider" that interacts with the effectController and also the rgbStripContoller (via effectController)
|
||||
# this could be seperated in one websocket server for the frontend and one for the rgbStrips
|
||||
# or an other frontend / rgbStrip backend provider not using websockets. you could integrate alexa, phillips hue like lamps or whatever you like!
|
||||
# but then there must be some autoloading of modules in a folder like the effects for easy installing. //todo :)
|
||||
print("starting websocket:8001")
|
||||
import BackendProvider.WebSocketServer as WebSocketServer
|
||||
webSocketThread = WebSocketServer.ThreadedWebSocketServer(
|
||||
effectController, rgbStripController)
|
||||
|
||||
print("starting UDPServer:8002")
|
||||
import BackendProvider.WemosStripUDPServer as UPDSocketServer
|
||||
udpSocketThread = UPDSocketServer.ThreadedUDPServer(
|
||||
effectController, rgbStripController)
|
||||
|
||||
while running:
|
||||
time.sleep(1)
|
||||
|
||||
except Exception as e:
|
||||
running = False
|
||||
print(e, traceback.format_exc())
|
||||
finally:
|
||||
print('shutting down the LED-Server')
|
||||
webSocketThread.stop()
|
||||
udpSocketThread.stop()
|
||||
effectController.stopAll()
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
main()
|
||||
@@ -0,0 +1,78 @@
|
||||
#!/usr/bin/env python
|
||||
# -*- coding: utf-8 -*-
|
||||
from os import path
|
||||
|
||||
#RGBStrip("Unter Theke",
|
||||
# wiringpi 24,-> BCM 19 -> GPIO. 24 -> kabel fehlt
|
||||
# wiringpi 4, -> BCM 23 -> GPIO. 4 -> Mosfet LAHMT (rot)
|
||||
# wiringpi 0, -> BCM 17 -> GPIO. 0 -> Mosfet TOT (weiß)
|
||||
#RGBStrip("Über Theke",
|
||||
# wiringpi 3, -> BCM 22 -> GPIO. 3
|
||||
# wiringpi 23,-> BCM 13 -> GPIO. 23
|
||||
# wiringpi 2, -> BCM 27 -> GPIO 2
|
||||
|
||||
#RGBStrip("Fensterbank",
|
||||
# wiringpi 21,-> BCM 5 -> GPIO. 21
|
||||
# wiringpi 25,-> BCM 26 -> GPIO. 25
|
||||
# wiringpi 22,-> BCM 6 -> GPIO. 22
|
||||
|
||||
#use the BCM pin numbers here
|
||||
"""
|
||||
from rgbUtils.RGBStrip import RGBStrip
|
||||
rgbStrips = [
|
||||
#RGBStrip("Test Dahem", 4, 17 , 22),
|
||||
RGBStrip("Unter Theke", 20, 16 , 21),
|
||||
RGBStrip("Über Theke", 22, 13, 27),
|
||||
RGBStrip("Fensterbank", 5, 26, 6)
|
||||
]
|
||||
|
||||
# setup PRi.GPIO
|
||||
GPIO.setmode(GPIO.BCM)
|
||||
# setup PWM for the rgbStrips
|
||||
for RGBStrip in self.getRGBStrips():
|
||||
RGBStrip.init()
|
||||
|
||||
"""
|
||||
"""
|
||||
Use WS2812B Strips:
|
||||
an arduino (uno tested) must be connected via usb while running
|
||||
the sketch in the root folder. Define the Strip as masterstrip and
|
||||
use parts of it as a rgbStrip
|
||||
|
||||
|
||||
"""
|
||||
|
||||
from rgbUtils.WS2812MasterStrip import WS2812MasterStrip
|
||||
from rgbUtils.WS2812Strip import WS2812Strip
|
||||
# LED_COUNT must be the same than in the arduino sketch
|
||||
ws2812master = WS2812MasterStrip('/dev/ttyACM0',150)
|
||||
rgbStrips = [
|
||||
WS2812Strip("LEDS 1-50",1,50,ws2812master),
|
||||
WS2812Strip("LEDS 51-100",51,100,ws2812master),
|
||||
WS2812Strip("LEDS 101-150",101,150,ws2812master),
|
||||
]
|
||||
|
||||
"""
|
||||
def drange(start, stop, step):
|
||||
r = start
|
||||
while r < stop:
|
||||
yield r
|
||||
r += step
|
||||
|
||||
rgbStrips = []
|
||||
for x in drange(1,150,10):
|
||||
rgbStrips.append(WS2812Strip(str(x)+"-"+str(x),x,x,ws2812master))
|
||||
"""
|
||||
|
||||
# int Port to bind. ports < 1024 need sudo access
|
||||
SocketBindPort = 8000
|
||||
|
||||
# Maximum brightness of the RGB Strips Max Value is 100, can be set lower if the strips are too bright.
|
||||
# (What I do not think, RGB Strips are never too bright)
|
||||
# MaxBrightness = 100
|
||||
|
||||
# GPIO Pins that are working with pwm. At the moment A and B models only
|
||||
# todo: check rpi version and add the missing pins if there are more that can be used
|
||||
AllowedGPIOPins = [3, 5, 7, 8, 10, 11, 12, 13, 15, 19, 21, 22, 23, 24, 26]
|
||||
|
||||
BASE_PATH = path.dirname(path.realpath(__file__))
|
||||
@@ -0,0 +1,240 @@
|
||||
from rgbUtils.BaseEffect import BaseEffect
|
||||
from rgbUtils.EffectParameter import slider, colorpicker
|
||||
from rgbUtils.debug import debug
|
||||
import time
|
||||
|
||||
import sys
|
||||
import numpy
|
||||
from rgbUtils.pyAudioRecorder import pyAudioRecorder
|
||||
from time import perf_counter, sleep
|
||||
from random import randint, shuffle
|
||||
|
||||
class musikEffect(BaseEffect):
|
||||
name = "musikEffect"
|
||||
desc = "LED-Band *sollte* nach musik blinken"
|
||||
|
||||
|
||||
|
||||
# Something that will be used to show descriptions and value options
|
||||
# of the parameters the effect will accept, in a way, that eg the webclient can decide,
|
||||
# if the parameters can be toggeled by a button/checkbox/slider/whatever
|
||||
effectParameters = [
|
||||
# radio(\
|
||||
# "Shuffle LED to Freq Order",\
|
||||
# "Off -> Mapping ",\
|
||||
# [\
|
||||
# [0,255,255,"red"],\
|
||||
# [0,255,255,"green"],\
|
||||
# [0,255,255,"blue"]\
|
||||
# ]\
|
||||
# ),\
|
||||
# slider(\
|
||||
# "Effect Brightnes",\
|
||||
# "Choose a brightness for your LED's",\
|
||||
# [\
|
||||
# [0,100,100,"brightness"],\
|
||||
# ]\
|
||||
# )\
|
||||
]
|
||||
|
||||
def init(self):
|
||||
|
||||
self.fft_random_keys = [0,1,2,3]
|
||||
self.fft_random = [0,0,0,0]
|
||||
|
||||
# used by strobe()
|
||||
self.lastmode = 0
|
||||
|
||||
self.recorderClient = pyAudioRecorder.recorderClient()
|
||||
self.rgbStripController.pyAudioRecorder.registerRecorderClient(self.recorderClient)
|
||||
|
||||
return
|
||||
|
||||
#loop effect as long as not stopped
|
||||
def effect(self):
|
||||
self.plot_audio_and_detect_beats()
|
||||
#self.freqtocolor()
|
||||
#if(time.time() - self.lastTime >= 0.002):
|
||||
#for RGBStrip in self.effectRGBStrips():
|
||||
# r= RGBStrip.red-1
|
||||
# if(r<0):
|
||||
# r=0
|
||||
# g= RGBStrip.green-1
|
||||
# if(g<0):
|
||||
# g=0
|
||||
# b= RGBStrip.blue-1
|
||||
# if(b<0):
|
||||
# b=0
|
||||
# RGBStrip.RGB(r,g,b)
|
||||
#self.lastTime = time.time()
|
||||
sleep(.001)
|
||||
|
||||
def end(self):
|
||||
self.rgbStripController.pyAudioRecorder.unregisterRecorderClient(self.recorderClient)
|
||||
|
||||
def plot_audio_and_detect_beats(self):
|
||||
if not self.rgbStripController.pyAudioRecorder.has_new_audio:
|
||||
return
|
||||
|
||||
# get x and y values from FFT
|
||||
xs, ys = self.rgbStripController.pyAudioRecorder.fft()
|
||||
|
||||
# calculate average for all frequency ranges
|
||||
y_avg = numpy.mean(ys)
|
||||
|
||||
|
||||
|
||||
#low_freq = numpy.mean(ys[20:47])
|
||||
#mid_freq = numpy.mean(ys[88:115])
|
||||
#hig_freq = numpy.mean(ys[156:184])
|
||||
|
||||
low_freq = numpy.mean(ys[0:67])
|
||||
mid_freq = numpy.mean(ys[68:135])
|
||||
hig_freq = numpy.mean(ys[136:204])
|
||||
|
||||
#get the maximum of all freq
|
||||
if len(self.y_max_freq_avg_list) < 250 or y_avg > 10 and numpy.amax([low_freq,mid_freq,hig_freq])/2 > numpy.amin(self.y_max_freq_avg_list):
|
||||
self.y_max_freq_avg_list.append(numpy.amax([low_freq,mid_freq,hig_freq]))
|
||||
|
||||
y_max = numpy.amax(self.y_max_freq_avg_list)
|
||||
|
||||
#y_max = numpy.mean([numpy.amax(ys),y_max])
|
||||
#print(low_freq,mid_freq,hig_freq,y_max)
|
||||
|
||||
for i,item in enumerate([low_freq,mid_freq,hig_freq]):
|
||||
if item is None:
|
||||
item = 0
|
||||
|
||||
low = round(low_freq/(y_max+10)*255)
|
||||
mid = round(mid_freq/(y_max+10)*255)
|
||||
hig = round(hig_freq/(y_max+10)*255)
|
||||
#print(low,mid,hig,y_max, numpy.amax(ys), y_avg)
|
||||
#print("------")
|
||||
|
||||
self.fft_random[self.fft_random_keys[0]] = low
|
||||
self.fft_random[self.fft_random_keys[1]] = mid
|
||||
self.fft_random[self.fft_random_keys[2]] = hig
|
||||
self.fft_random[self.fft_random_keys[3]] = 0
|
||||
|
||||
# calculate low frequency average
|
||||
#low_freq = [ys[i] for i in range(len(xs)) if xs[i] < 1000]
|
||||
low_freq = ys[0:47]
|
||||
low_freq_avg = numpy.mean(low_freq)
|
||||
|
||||
if len(self.low_freq_avg_list) < 250 or low_freq_avg > numpy.amin(self.low_freq_avg_list)/2:
|
||||
self.low_freq_avg_list.append(low_freq_avg)
|
||||
cumulative_avg = numpy.mean(self.low_freq_avg_list)
|
||||
|
||||
bass = low_freq[:int(len(low_freq)/2)]
|
||||
bass_avg = numpy.mean(bass)
|
||||
#print("bass: {:.2f} vs cumulative: {:.2f}".format(bass_avg, cumulative_avg))
|
||||
|
||||
# check if there is a beat
|
||||
# song is pretty uniform across all frequencies
|
||||
if (y_avg > y_avg/5 and (bass_avg > cumulative_avg * 1.8 or (low_freq_avg < y_avg * 1.2 and bass_avg > cumulative_avg))):
|
||||
#self.prev_beat
|
||||
curr_time = perf_counter()
|
||||
|
||||
# print(curr_time - self.prev_beat)
|
||||
if curr_time - self.prev_beat > 60/360*2: # 180 BPM max
|
||||
shuffle(self.fft_random_keys)
|
||||
|
||||
# change the button color
|
||||
#self.beats_idx += 1
|
||||
#self.strobe()
|
||||
#print("beat {}".format(self.beats_idx))
|
||||
#print("bass: {:.2f} vs cumulative: {:.2f}".format(bass_avg, cumulative_avg))
|
||||
|
||||
#print(self.fft_random)
|
||||
# change the button text
|
||||
bpm = int(60 / (curr_time - self.prev_beat))
|
||||
if len(self.bpm_list) < 4:
|
||||
if bpm > 60:
|
||||
self.bpm_list.append(bpm)
|
||||
else:
|
||||
bpm_avg = int(numpy.mean(self.bpm_list))
|
||||
if abs(bpm_avg - bpm) < 35:
|
||||
self.bpm_list.append(bpm)
|
||||
print("bpm: {:d}".format(bpm_avg))
|
||||
|
||||
# reset the timer
|
||||
self.prev_beat = curr_time
|
||||
if y_avg > 10:
|
||||
for RGBStrip in self.effectRGBStrips():
|
||||
RGBStrip.RGB(
|
||||
self.fft_random[0],
|
||||
self.fft_random[1],
|
||||
self.fft_random[2]
|
||||
)
|
||||
|
||||
# shorten the cumulative list to account for changes in dynamics
|
||||
if len(self.low_freq_avg_list) > 500:
|
||||
self.low_freq_avg_list = self.low_freq_avg_list[250:]
|
||||
#print("REFRESH!!")
|
||||
|
||||
# shorten the cumulative list to account for changes in dynamics
|
||||
if len(self.y_max_freq_avg_list ) > 500:
|
||||
self.y_max_freq_avg_list = self.y_max_freq_avg_list[250:]
|
||||
print("--REFRESH y_max_freq_avg_list")
|
||||
|
||||
# keep two 8-counts of BPMs so we can maybe catch tempo changes
|
||||
if len(self.bpm_list) > 24:
|
||||
self.bpm_list = self.bpm_list[8:]
|
||||
|
||||
# reset song data if the song has stopped
|
||||
if y_avg < 10:
|
||||
self.bpm_list = []
|
||||
self.low_freq_avg_list = []
|
||||
print("new song")
|
||||
self.off()
|
||||
|
||||
self.rgbStripController.pyAudioRecorder.newAudio = False
|
||||
# print(self.bpm_list)250
|
||||
|
||||
def strobe(self):
|
||||
x = randint(0,5)
|
||||
while x is self.lastmode:
|
||||
x = randint(0,5)
|
||||
|
||||
self.lastmode = x
|
||||
r = 255#randint(0,255)
|
||||
g = 255#randint(0,255)
|
||||
b = 255#randint(0,255)
|
||||
if x is 0:
|
||||
for RGBStrip in self.effectRGBStrips():
|
||||
RGBStrip.RGB(r,g,0)
|
||||
if x is 1:
|
||||
for RGBStrip in self.effectRGBStrips():
|
||||
RGBStrip.RGB(0,g,b)
|
||||
if x is 2:
|
||||
for RGBStrip in self.effectRGBStrips():
|
||||
RGBStrip.RGB(r,0,b)
|
||||
if x is 3:
|
||||
for RGBStrip in self.effectRGBStrips():
|
||||
RGBStrip.RGB(r,0,0)
|
||||
if x is 4:
|
||||
for RGBStrip in self.effectRGBStrips():
|
||||
RGBStrip.RGB(0,g,0)
|
||||
if x is 5:
|
||||
for RGBStrip in self.effectRGBStrips():
|
||||
RGBStrip.RGB(0,0,b)
|
||||
if x is 6:
|
||||
for RGBStrip in self.effectRGBStrips():
|
||||
RGBStrip.RGB(r,g,b)
|
||||
|
||||
def off(self):
|
||||
for RGBStrip in self.effectRGBStrips():
|
||||
RGBStrip.RGB(0,0,0)
|
||||
|
||||
def left_rotate(self,arr):
|
||||
if not arr:
|
||||
return arr
|
||||
|
||||
left_most_element = arr[0]
|
||||
length = len(arr)
|
||||
|
||||
for i in range(length - 1):
|
||||
arr[i], arr[i + 1] = arr[i + 1], arr[i]
|
||||
|
||||
arr[length - 1] = left_most_element
|
||||
return arr
|
||||
@@ -0,0 +1,23 @@
|
||||
from rgbUtils.BaseEffect import BaseEffect
|
||||
from rgbUtils.debug import debug
|
||||
|
||||
import time
|
||||
|
||||
class offEffect(BaseEffect):
|
||||
name = "offEffect"
|
||||
desc = "LED-Band *sollte* nicht an sein"
|
||||
|
||||
def effect(self):
|
||||
time.sleep(1)
|
||||
return
|
||||
|
||||
# for overriding by the effect, when a strip is added
|
||||
def onRGBStripAdded(self,rgbStrip):
|
||||
rgbStrip.RGB(0,0,0,)
|
||||
return
|
||||
|
||||
# for overriding by the effect, when a strip is added
|
||||
def onEffectParameterValuesUpdated(self):
|
||||
for RGBStrip in self.effectRGBStrips():
|
||||
RGBStrip.RGB(0,0,0)
|
||||
return
|
||||
@@ -0,0 +1,67 @@
|
||||
from rgbUtils.BaseEffect import BaseEffect
|
||||
from rgbUtils.EffectParameter import slider, colorpicker
|
||||
from rgbUtils.debug import debug
|
||||
import time
|
||||
|
||||
class onEffect(BaseEffect):
|
||||
name = "onEffect"
|
||||
desc = "LED-Band *sollte* an sein"
|
||||
|
||||
# Something that will be used to show descriptions and value options
|
||||
# of the parameters the effect will accept, in a way, that eg the webclient can decide,
|
||||
# if the parameters can be toggeled by a button/checkbox/slider/whatever
|
||||
effectParameters = [
|
||||
colorpicker(\
|
||||
"Effect Color",\
|
||||
"Choose a color for your LED's",\
|
||||
[\
|
||||
[0,255,255,"red"],\
|
||||
[0,255,255,"green"],\
|
||||
[0,255,255,"blue"]\
|
||||
]\
|
||||
),\
|
||||
slider(\
|
||||
"Effect Brightnes",\
|
||||
"Choose a brightness for your LED's",\
|
||||
[\
|
||||
[0,100,100,"brightness"],\
|
||||
]\
|
||||
)\
|
||||
]
|
||||
|
||||
def init(self):
|
||||
return
|
||||
|
||||
#loop effect as long as not stopped
|
||||
def effect(self):
|
||||
time.sleep(1)
|
||||
return
|
||||
|
||||
# for overriding by the effect, when a strip is added
|
||||
def onRGBStripAdded(self,rgbStrip):
|
||||
rgbStrip.RGB(\
|
||||
# colorpicker red currentvalue
|
||||
self.effectParameterValues[0][0],\
|
||||
# colorpicker green currentvalue
|
||||
self.effectParameterValues[0][1],\
|
||||
# colorpicker blue currentvalue
|
||||
self.effectParameterValues[0][2],\
|
||||
# slider brightness currentvalue
|
||||
self.effectParameterValues[1][0]\
|
||||
)
|
||||
return
|
||||
|
||||
# for overriding by the effect, when a params are updated
|
||||
def onEffectParameterValuesUpdated(self):
|
||||
for RGBStrip in self.effectRGBStrips():
|
||||
#print(self.effectParameterValues)
|
||||
RGBStrip.RGB(\
|
||||
# colorpicker red currentvalue
|
||||
self.effectParameterValues[0][0],\
|
||||
# colorpicker green currentvalue
|
||||
self.effectParameterValues[0][1],\
|
||||
# colorpicker blue currentvalue
|
||||
self.effectParameterValues[0][2],\
|
||||
# slider brightness currentvalue
|
||||
self.effectParameterValues[1][0]\
|
||||
)
|
||||
@@ -0,0 +1,67 @@
|
||||
from rgbUtils.BaseEffect import BaseEffect
|
||||
from rgbUtils.debug import debug
|
||||
import time
|
||||
|
||||
class rainbowEffect(BaseEffect):
|
||||
name = "rainbowEffect"
|
||||
desc = "LED-Band *sollte* rainbowEffect sein"
|
||||
|
||||
def init(self):
|
||||
self.i=0
|
||||
self.speed = 1
|
||||
self.helligkeit = 100
|
||||
|
||||
#loop effect as long as not stopped
|
||||
def effect(self):
|
||||
debug(self)
|
||||
if self.i < 3*255*self.speed:
|
||||
c = self.wheel_color(self.i,self.speed)
|
||||
#print("r: "+ str(round(c[0]/100*helligkeit,2))+" g: "+str(round(c[1]/100*helligkeit,2))+" b: "+str(round(c[2]/100*helligkeit,2)))
|
||||
for rgbStrip in self.effectRGBStrips():
|
||||
#print(c[0],c[1],c[2])
|
||||
rgbStrip.RGB(c[0],c[1],c[2],self.helligkeit)
|
||||
time.sleep(0.01)
|
||||
self.i =self.i +1
|
||||
else:
|
||||
self.i=0
|
||||
|
||||
# for overriding by the effect, when a strip is added
|
||||
def onRGBStripAdded(self,rgbStrip):
|
||||
return
|
||||
|
||||
# for overriding by the effect, when a strip is added
|
||||
def onEffectParameterValuesUpdated(self):
|
||||
return
|
||||
|
||||
def wheel_color(self,position,speed = 5):
|
||||
"""Get color from wheel value (0 - 765)"""
|
||||
if position < 0:
|
||||
position = 0
|
||||
if position > 765*speed:
|
||||
position = 765*speed
|
||||
|
||||
if position < (255*speed):
|
||||
r = (255*speed) - position % (255*speed)
|
||||
g = position % (255*speed)
|
||||
b = 0
|
||||
elif position < (510*speed):
|
||||
g = (255*speed) - position % (255*speed)
|
||||
b = position % (255*speed)
|
||||
r = 0
|
||||
else:
|
||||
b = (255*speed) - position % (255*speed)
|
||||
r = position % (255*speed)
|
||||
g = 0
|
||||
|
||||
return [r/speed, g/speed, b/speed]
|
||||
|
||||
# def wheel_color_2(self,r=255,g=0,b=0):
|
||||
# if r<255:
|
||||
# r=r+1
|
||||
# elif g<255:
|
||||
# g=g+1
|
||||
# elif r=255:
|
||||
|
||||
# elif b<255:
|
||||
# b=b+1
|
||||
|
||||
@@ -0,0 +1,42 @@
|
||||
from rgbUtils.BaseEffect import BaseEffect
|
||||
from rgbUtils.debug import debug
|
||||
from random import randint
|
||||
import time
|
||||
|
||||
class strobeEffect(BaseEffect):
|
||||
name = "strobeEffect"
|
||||
desc = "*Strobe*"
|
||||
|
||||
def init(self):
|
||||
self.state = True
|
||||
|
||||
#loop effect as long as not stopped
|
||||
def effect(self):
|
||||
y = -1
|
||||
x = -1
|
||||
if self.state:
|
||||
while x is y:
|
||||
x = randint(0,2)
|
||||
if x is 0:
|
||||
for RGBStrip in self.effectRGBStrips():
|
||||
RGBStrip.RGB(255,255,0)
|
||||
if x is 1:
|
||||
for RGBStrip in self.effectRGBStrips():
|
||||
RGBStrip.RGB(0,255,255)
|
||||
if x is 2:
|
||||
for RGBStrip in self.effectRGBStrips():
|
||||
RGBStrip.RGB(255,0,255)
|
||||
self.state = False
|
||||
else:
|
||||
for RGBStrip in self.effectRGBStrips():
|
||||
RGBStrip.RGB(0,0,0)
|
||||
self.state = True
|
||||
time.sleep(0.01)
|
||||
|
||||
# for overriding by the effect, when a strip is added
|
||||
def onRGBStripAdded(self,rgbStrip):
|
||||
return
|
||||
|
||||
# for overriding by the effect, when a strip is added
|
||||
def onEffectParameterValuesUpdated(self):
|
||||
return
|
||||
@@ -0,0 +1,91 @@
|
||||
import time
|
||||
import threading
|
||||
import copy
|
||||
from rgbUtils.debug import debug
|
||||
|
||||
class BaseEffect(threading.Thread):
|
||||
|
||||
# The Name and the Description of the Effect,
|
||||
# should be overwritten by the inheritancing Effect
|
||||
name = "Undefined"
|
||||
desc = "No Description"
|
||||
|
||||
# Something that will be used to show descriptions and value options
|
||||
# of the parameters the effect will accept, in a way, that eg the webclient can decide,
|
||||
# if the parameters can be toggeled by a button/checkbox/slider/whatever
|
||||
effectParameters = []
|
||||
|
||||
stop = False
|
||||
|
||||
def __init__(self):
|
||||
threading.Thread.__init__(self)
|
||||
self.effectRGBStripList = []
|
||||
|
||||
# when a strip is added or removed or the options change, the thread will not restart,
|
||||
# the changes are pushed to the thread at livetime. so avoid log running loops without
|
||||
# accessing getRGBStrips and getEffectParams. It could happen, that a new effect for example already uses
|
||||
# a rgbStip while the loop of the old effect is running and has not got the changes.
|
||||
def run(self):
|
||||
self.effectParameterValues = []
|
||||
for effectParameterIndex, effectParameter in enumerate(self.effectParameters):
|
||||
self.effectParameterValues.append([])
|
||||
for option in effectParameter.options:
|
||||
self.effectParameterValues[effectParameterIndex].append(option[2])
|
||||
|
||||
print(self.effectParameterValues)
|
||||
self.init()
|
||||
while not self.stop:
|
||||
# run the effect in endless while
|
||||
self.effect()
|
||||
self.end()
|
||||
|
||||
# Init is called bevor the loop, use setEffectParams for init values
|
||||
# and access them in the effect with getEffectParams
|
||||
def init(self):
|
||||
return
|
||||
|
||||
# see run(): never save effectRGBStrips() as a variable and access it as often as possible
|
||||
# to avoid two effects accessing the rgbStrip
|
||||
def effect(self):
|
||||
while 1:
|
||||
debug("ET "+self.name+" effect() function needs to be replaced in inheritancing Effect")
|
||||
|
||||
# called when the effect is stopped
|
||||
def end(self):
|
||||
return
|
||||
|
||||
# when called the effect loop will stop and the thread can be terminated
|
||||
def stopEffect(self):
|
||||
self.stop = True
|
||||
|
||||
# the effect itself will know its own strips. i don't know if it would be better if the effectController self
|
||||
# should have this list, but then i must do something like nested arrays, naah.
|
||||
def addRGBStrip(self,rgbStrip):
|
||||
self.effectRGBStripList.append(rgbStrip)
|
||||
self.onRGBStripAdded(rgbStrip)
|
||||
|
||||
# remove a strip, if the effect has no more strips, the effect thread guardian will kill it
|
||||
def removeRGBStrip(self,rgbStrip):
|
||||
self.effectRGBStripList.remove(rgbStrip)
|
||||
|
||||
# for overriding by the effect, when a strip is added
|
||||
def onRGBStripAdded(self,rgbStrip):
|
||||
return
|
||||
# for overriding by the effect, when a strip is added
|
||||
def onEffectParameterValuesUpdated(self):
|
||||
return
|
||||
|
||||
# returns a list of the RGBStrips used by this effect
|
||||
def effectRGBStrips(self):
|
||||
return self.effectRGBStripList
|
||||
|
||||
def addEffectParameter(self,effectParameter):
|
||||
self.effectParameters.append(effectParameter)
|
||||
|
||||
# set Params as descriped in getParamsDescription()
|
||||
def updateEffectParameterValues(self,effectParameterIndex, effectParameterValues):
|
||||
print("updateEffectParameterValues",effectParameterIndex,effectParameterValues)
|
||||
for effectParameterValue in effectParameterValues:
|
||||
if self.effectParameters[int(effectParameterIndex)].testValue(int(effectParameterValue),int(effectParameterValues[effectParameterValue])):
|
||||
self.effectParameterValues[int(effectParameterIndex)][int(effectParameterValue)] = int(effectParameterValues[effectParameterValue])
|
||||
self.onEffectParameterValuesUpdated()
|
||||
@@ -0,0 +1,42 @@
|
||||
import threading
|
||||
|
||||
class effectParameter(object):
|
||||
# The Name and the Description of the EffectParameter,
|
||||
# should be overwritten by the inheritancing EffectParameter
|
||||
name="Undefined"
|
||||
desc = "No Description"
|
||||
|
||||
# In the order you expect the options to be set
|
||||
# [
|
||||
# [min/off,max/on,current,"description"],
|
||||
# [min/off,max/on,current,"description"],
|
||||
# [min/off,max/on,current,"description"],
|
||||
# ]
|
||||
options = []
|
||||
|
||||
def __init__(self,name,desc,initOptions = []):
|
||||
self.name = name
|
||||
self.desc = desc
|
||||
self.options = initOptions
|
||||
|
||||
# check if the given values are plausible
|
||||
def testValue(self,index,value):
|
||||
if value >= self.options[index][0] \
|
||||
and value <= self.options[index][1]:
|
||||
return True
|
||||
else:
|
||||
return False
|
||||
|
||||
class colorpicker(effectParameter):
|
||||
name="UndefinedColorpicker"
|
||||
desc="No Description"
|
||||
type="colorpicker"
|
||||
|
||||
# check if the given values are plausible
|
||||
def testValue(self,index,value):
|
||||
return True
|
||||
|
||||
class slider(effectParameter):
|
||||
name="UndefinedSlider"
|
||||
desc="No Description"
|
||||
type="slider"
|
||||
@@ -0,0 +1,95 @@
|
||||
from rgbUtils.debug import debug
|
||||
import uuid
|
||||
|
||||
class RGBStrip:
|
||||
# name = the name off the the strip, defined by the client connecting to the server
|
||||
# uid = unique id, if the strip sends one, use this (later maybe, or never, whatever)
|
||||
# lenght = the lenght off the strip, for future use of eg WS2812b strips, will be 1 by default
|
||||
def __init__(self,name,onValuesUpdateHandler,lenght=1):
|
||||
# UID should be updateable later, or not?
|
||||
# when updating, be sure it does not exist
|
||||
self.STRIP_UID = str(uuid.uuid4())
|
||||
self.STRIP_NAME = name
|
||||
self.STRIP_LEGHT = lenght
|
||||
|
||||
self.onValuesUpdateHandler = onValuesUpdateHandler
|
||||
|
||||
self.red = [0]*self.STRIP_LEGHT
|
||||
self.green = [0]*self.STRIP_LEGHT
|
||||
self.blue = [0]*self.STRIP_LEGHT
|
||||
|
||||
def RGB(self,red,green,blue,brightness = 100):
|
||||
|
||||
if(red < 0):
|
||||
red = 0
|
||||
if(red > 255):
|
||||
red = 255
|
||||
|
||||
if(green < 0):
|
||||
green = 0
|
||||
if(green > 255):
|
||||
green = 255
|
||||
|
||||
if(blue < 0):
|
||||
blue = 0
|
||||
if(blue > 255):
|
||||
blue = 255
|
||||
|
||||
if(brightness < 0):
|
||||
brightness = 0
|
||||
if(brightness > 100):
|
||||
brightness = 100
|
||||
|
||||
for x in range(self.STRIP_LEGHT):
|
||||
self.red[x] = int(red/100*brightness)
|
||||
self.green[x] = int(green/100*brightness)
|
||||
self.blue[x] = int(blue/100*brightness)
|
||||
|
||||
self.onValuesUpdateHandler(self)
|
||||
|
||||
def WS2812b(self,id,red,green,blue,brightness=100):
|
||||
if id < 0 and id > self.STRIP_LEGHT:
|
||||
print(self.STRIP_NAME," is max ",self.STRIP_LEGHT," Pixels long!")
|
||||
return
|
||||
else:
|
||||
if(red < 0):
|
||||
red = 0
|
||||
if(red > 255):
|
||||
red = 255
|
||||
|
||||
if(green < 0):
|
||||
green = 0
|
||||
if(green > 255):
|
||||
green = 255
|
||||
|
||||
if(blue < 0):
|
||||
blue = 0
|
||||
if(blue > 255):
|
||||
blue = 255
|
||||
|
||||
if(brightness < 0):
|
||||
brightness = 0
|
||||
if(brightness > 100):
|
||||
brightness = 100
|
||||
|
||||
self.red[id] = int(red/100*brightness)
|
||||
self.green[id] = int(green/100*brightness)
|
||||
self.blue[id] = int(blue/100*brightness)
|
||||
|
||||
self.onValuesUpdateHandler(self,id)
|
||||
|
||||
|
||||
def off(self):
|
||||
for x in range(self.STRIP_LEGHT):
|
||||
self.red[x] = 0
|
||||
self.green[x] = 0
|
||||
self.blue[x] = 0
|
||||
|
||||
self.onValuesUpdateHandler(self)
|
||||
|
||||
def getData(self):
|
||||
self.hasNewData = False
|
||||
return [self.red,self.green,self.blue]
|
||||
|
||||
|
||||
|
||||
@@ -0,0 +1,6 @@
|
||||
#!/usr/bin/env python
|
||||
# -*- coding: utf-8 -*-
|
||||
|
||||
def debug(string):
|
||||
return
|
||||
print(string)
|
||||
@@ -0,0 +1,179 @@
|
||||
# import offEffect
|
||||
from effects.offEffect import offEffect
|
||||
|
||||
from rgbUtils.debug import debug
|
||||
|
||||
from rgbUtils.BaseEffect import BaseEffect
|
||||
from rgbUtils.RGBStrip import RGBStrip
|
||||
|
||||
import time
|
||||
import threading
|
||||
import os
|
||||
import sys
|
||||
|
||||
#
|
||||
# handles all the effects:
|
||||
# - start/stop effects
|
||||
# - changes parameters of specific effects
|
||||
# - moves the rgbStrips to the effects
|
||||
# - runs a guardian that detects dead effectThreads and moves Strips to offEffect
|
||||
#
|
||||
class effectController:
|
||||
|
||||
# list of current running effects
|
||||
effectThreads = []
|
||||
|
||||
offEffectThreadObject = None
|
||||
|
||||
# maybe i will seperate this later
|
||||
onControllerChangeHandler = []
|
||||
|
||||
def __init__(self,rgbStripController):
|
||||
self.rgbStripController = rgbStripController
|
||||
# load the effects
|
||||
self.effectsList = self.getEffectsListFromDir()
|
||||
# start the offEffect by default
|
||||
self.offEffectThreadObject = self.startEffect(offEffect,[])
|
||||
# - a bit of failover handling, remove dead threads from effectThread array
|
||||
# - move strips without an effect to the offEffect
|
||||
self.effectGuardian = self.effectGuardian(self)
|
||||
self.effectGuardian.start()
|
||||
|
||||
# starts a effect by given class, rgbStrips and params ([[index,[param,param,param]],[index,[param,param,param]]])
|
||||
def startEffect(self, effectClass: BaseEffect, rgbStrips: list, params: list = []):
|
||||
|
||||
newEffect = effectClass()
|
||||
if len(self.effectThreads) > 0 and len(rgbStrips) > 0 or len(self.effectThreads) == 0:
|
||||
newEffect.start()
|
||||
self.updateEffectParameters(newEffect, params)
|
||||
self.effectThreads.append(newEffect)
|
||||
|
||||
for rgbStrip in rgbStrips:
|
||||
self.moveRGBStripToEffectThread(rgbStrip, newEffect)
|
||||
|
||||
self.noticeControllerChange()
|
||||
return newEffect
|
||||
|
||||
# returns all effectClasses but offEffect, since offEffect will never be killed
|
||||
# and should not be running twice
|
||||
def getEffects(self):
|
||||
# alle außer offEffect
|
||||
return self.effectsList
|
||||
|
||||
# returns all running effectThreads
|
||||
def getEffectThreads(self):
|
||||
return self.effectThreads
|
||||
|
||||
# returns a list of the RGBStrips used by this effect
|
||||
def getEffectRGBStrips(self, effectThreadObject: BaseEffect):
|
||||
return effectThreadObject.effectRGBStrips()
|
||||
|
||||
# move a rgbStip to a running effectThread
|
||||
def moveRGBStripToEffectThread(self, rgbStrip: RGBStrip, effectThreadObject: BaseEffect):
|
||||
# cycle throught all effects and
|
||||
# remove Strip from effect if added
|
||||
for et in self.effectThreads:
|
||||
if rgbStrip in et.effectRGBStrips():
|
||||
et.removeRGBStrip(rgbStrip)
|
||||
if effectThreadObject.isAlive():
|
||||
effectThreadObject.addRGBStrip(rgbStrip)
|
||||
# check if any effectThread has no more rgbStrips and if so, stop it
|
||||
|
||||
# if the effectThread has no more strips, we stop it and remove it.
|
||||
for x, effectThread in enumerate(self.effectThreads):
|
||||
if len(effectThread.effectRGBStrips()) == 0 and x is not 0:
|
||||
effectThread.stopEffect()
|
||||
self.effectThreads.remove(effectThread)
|
||||
self.noticeControllerChange()
|
||||
|
||||
# updates parameter of a running effectThread
|
||||
def updateEffectParameters(self, effectThreadObject: BaseEffect, effectParameters):
|
||||
for effectParameter in effectParameters:
|
||||
effectThreadObject.updateEffectParameterValues(
|
||||
effectParameter[0], effectParameter[1])
|
||||
self.noticeControllerChange()
|
||||
|
||||
# stops all effectThreads and set the rgbStrips to off
|
||||
def stopAll(self):
|
||||
debug("effectController stopAll()")
|
||||
for effectThread in self.effectThreads:
|
||||
effectThread.stopEffect()
|
||||
debug("effectController killed "+str(effectThread))
|
||||
|
||||
self.effectGuardian.stop()
|
||||
|
||||
time.sleep(0.5)
|
||||
# GPIO.cleanup()
|
||||
|
||||
# inform the controllerChangeHandler to update the client
|
||||
def noticeControllerChange(self):
|
||||
for controllerChangeHandler in self.onControllerChangeHandler:
|
||||
controllerChangeHandler()
|
||||
|
||||
# add onControllerChangeHandler
|
||||
def addOnControllerChangeHandler(self, hander):
|
||||
print("addOnControllerChangeHandler", str(hander))
|
||||
self.onControllerChangeHandler.append(hander)
|
||||
# send data to this client
|
||||
hander()
|
||||
|
||||
# remove onControllerChangeHandler
|
||||
def removeOnControllerChangeHandler(self, hander):
|
||||
print("removeOnControllerChangeHandler", str(hander))
|
||||
if hander in self.onControllerChangeHandler:
|
||||
self.onControllerChangeHandler.remove(hander)
|
||||
else:
|
||||
print('\n\n -> client was never registered!')
|
||||
|
||||
# automaticly loads all modules from effects subdir and adds them to the list of effects if they have the BaseEffect as subclass
|
||||
def getEffectsListFromDir(self):
|
||||
effectsList = []
|
||||
for raw_module in os.listdir(os.path.abspath(os.path.join(os.path.join(os.path.dirname(__file__), os.pardir),"effects"))):
|
||||
|
||||
if raw_module == '__init__.py' or raw_module == 'offEffect.py' or raw_module[-3:] != '.py':
|
||||
continue
|
||||
effectModule = __import__("effects."+raw_module[:-3], fromlist=[raw_module[:-3]])
|
||||
effectClass = getattr(effectModule, raw_module[:-3])
|
||||
if issubclass(effectClass,BaseEffect):
|
||||
effectsList.append(effectClass)
|
||||
print("Loaded effects: ",effectsList)
|
||||
return effectsList
|
||||
|
||||
def onRGBStripRegistered(self,rgbStrip):
|
||||
self.offEffectThreadObject.addRGBStrip(rgbStrip)
|
||||
self.noticeControllerChange()
|
||||
|
||||
def onRGBStripUnRegistered(self,rgbStrip):
|
||||
# cycle throught all effects and
|
||||
# remove Strip from effect if added
|
||||
for et in self.effectThreads:
|
||||
if rgbStrip in et.effectRGBStrips():
|
||||
et.removeRGBStrip(rgbStrip)
|
||||
# if the effectThread has no more strips, we stop it and remove it.
|
||||
for x, effectThread in enumerate(self.effectThreads):
|
||||
if len(effectThread.effectRGBStrips()) == 0 and x is not 0:
|
||||
effectThread.stopEffect()
|
||||
self.effectThreads.remove(effectThread)
|
||||
self.noticeControllerChange()
|
||||
|
||||
class effectGuardian(threading.Thread):
|
||||
def __init__(self, effectController):
|
||||
threading.Thread.__init__(self)
|
||||
self.effectController = effectController
|
||||
self.stopped = False
|
||||
|
||||
def run(self):
|
||||
while not self.stopped:
|
||||
for effectThread in self.effectController.effectThreads:
|
||||
# if Thread was killed by something else, we remove it from the list
|
||||
if not effectThread.isAlive():
|
||||
for rgbStrip in effectThread.effectRGBStrips():
|
||||
self.effectController.moveRGBStripToEffectThread(rgbStrip,self.effectController.offEffectThreadObject)
|
||||
self.effectController.effectThreads.remove(effectThread)
|
||||
print("effectController:effectGuardian removed dead Thread " +
|
||||
str(effectThread) + ". There must be an error in code!")
|
||||
self.effectController.noticeControllerChange()
|
||||
time.sleep(1)
|
||||
|
||||
def stop(self):
|
||||
self.stopped = True
|
||||
@@ -0,0 +1,49 @@
|
||||
"""
|
||||
Convert the effectController function outputs to json format
|
||||
"""
|
||||
|
||||
def responseHandler(effectController,rgbStripController,data):
|
||||
if "startEffect" in data:
|
||||
enabledRGBStrips = []
|
||||
for rgbStrip in rgbStripController.getRGBStrips():
|
||||
for rgbStripJsonArray in data['startEffect']['rgbStrips']:
|
||||
if rgbStrip.STRIP_UID in rgbStripJsonArray[0] and rgbStripJsonArray[1]:
|
||||
enabledRGBStrips.append(rgbStrip)
|
||||
|
||||
effectController.startEffect(effectController.getEffects()[data['startEffect']['effect']],enabledRGBStrips,data['startEffect']['params'])
|
||||
if "moveRGBStripToEffectThread" in data:
|
||||
for rgbStrip in rgbStripController.getRGBStrips():
|
||||
if rgbStrip.STRIP_UID in data['moveRGBStripToEffectThread']['rgbStrip']:
|
||||
effectController.moveRGBStripToEffectThread( \
|
||||
rgbStrip, \
|
||||
effectController.getEffectThreads()[data['moveRGBStripToEffectThread']['effectThread']] \
|
||||
)
|
||||
if "effectThreadChangeEffectParam" in data:
|
||||
effectController.updateEffectParameters(\
|
||||
effectController.getEffectThreads()[data['effectThreadChangeEffectParam']['effectThread']], \
|
||||
data['effectThreadChangeEffectParam']['params']
|
||||
)
|
||||
return 'ok'
|
||||
|
||||
# return json of all configured effects (except offEffect) with their paramDescriptions
|
||||
def getEffects(effectController):
|
||||
result = {}
|
||||
for x, effect in enumerate(effectController.getEffects()):
|
||||
effectParams = {}
|
||||
for y, effectParam in enumerate(effect.effectParameters):
|
||||
effectParams[y] = {'index': y,'type': effectParam.type, 'name': effectParam.name, 'desc': effectParam.desc, 'options': effectParam.options}
|
||||
result[x] = {'index': x, 'name': effect.name, 'desc': effect.desc, 'effectParams': effectParams}
|
||||
return result
|
||||
|
||||
# return json of all running effectThreads with their active rgbStrips and params
|
||||
def getEffectThreads(effectController):
|
||||
result = {}
|
||||
for x, effectThread in enumerate(effectController.getEffectThreads()):
|
||||
effectRGBStrips = {}
|
||||
for effectRGBStrip in effectController.getEffectRGBStrips(effectThread):
|
||||
effectRGBStrips[effectRGBStrip.STRIP_UID] = {'index': effectRGBStrip.STRIP_UID, 'name': effectRGBStrip.STRIP_NAME}
|
||||
effectParams = {}
|
||||
for z, effectParam in enumerate(effectThread.effectParameters):
|
||||
effectParams[z] = {'index': z,'type': effectParam.type, 'name': effectParam.name, 'desc': effectParam.desc, 'options': effectParam.options, 'values': effectThread.effectParameterValues[z]}
|
||||
result[x] = {'index': x, 'name': effectThread.name, 'desc': effectThread.desc,'activeRGBStips': effectRGBStrips, 'dump': str(effectThread), 'effectParams': effectParams}
|
||||
return result
|
||||
@@ -0,0 +1,140 @@
|
||||
#
|
||||
# The idea is to have only one thread accessing the audio input source instead
|
||||
# of every music enabled thread itself. also, different fuctions calculating
|
||||
# frequencys and so on shoud run in own threads to the leds get more updates
|
||||
#
|
||||
# in history one thread was calculating bpm, freqence average and max values in one thread
|
||||
# before updating the leds, what the pi was a bit slow for and you cloud count every update
|
||||
# of the leds
|
||||
|
||||
# i think most of it is from https://github.com/shunfu/python-beat-detector/
|
||||
|
||||
|
||||
import numpy
|
||||
import scipy
|
||||
import pyaudio
|
||||
import threading
|
||||
|
||||
import time
|
||||
|
||||
|
||||
class pyAudioRecorder:
|
||||
"""Simple, cross-platform class to record from the default input device."""
|
||||
|
||||
def __init__(self):
|
||||
self.RATE = 44100
|
||||
self.BUFFERSIZE = 2**12
|
||||
self.secToRecord = .1
|
||||
self.kill_threads = False
|
||||
self.has_new_audio = False
|
||||
self.setup()
|
||||
|
||||
# since the server only can handly one input (for now) this thread will calculate
|
||||
# some basic things for the musicEffectsThreads, so they can update the leds more often.
|
||||
# it is always posible to catch the fft() function and do your own thing in your musikEffect.
|
||||
|
||||
# calculate bpm, this is the same for all clientThreads
|
||||
self.beats_idx = 0
|
||||
|
||||
self.bpm_list = []
|
||||
self.prev_beat = time.perf_counter()
|
||||
self.low_freq_avg_list = []
|
||||
|
||||
self.lastTime = time.time()
|
||||
|
||||
|
||||
self.recorderClients = []
|
||||
|
||||
def setup(self):
|
||||
self.buffers_to_record = int(
|
||||
self.RATE * self.secToRecord / self.BUFFERSIZE)
|
||||
if self.buffers_to_record == 0:
|
||||
self.buffers_to_record = 1
|
||||
self.samples_to_record = int(self.BUFFERSIZE * self.buffers_to_record)
|
||||
self.chunks_to_record = int(self.samples_to_record / self.BUFFERSIZE)
|
||||
self.sec_per_point = 1. / self.RATE
|
||||
|
||||
self.p = pyaudio.PyAudio()
|
||||
# start the PyAudio class
|
||||
for i in range(self.p.get_device_count()):
|
||||
devinfo = self.p.get_device_info_by_index(i)
|
||||
print(i, devinfo["name"])
|
||||
# make sure the default input device is broadcasting the speaker output
|
||||
# there are a few ways to do this
|
||||
# e.g., stereo mix, VB audio cable for windows, soundflower for mac
|
||||
self.in_stream = self.p.open(format=pyaudio.paInt16,
|
||||
channels=1,
|
||||
rate=self.RATE,
|
||||
input=True,
|
||||
frames_per_buffer=self.BUFFERSIZE)
|
||||
print("Using default input device: {:s}".format(
|
||||
self.p.get_default_input_device_info()['name']))
|
||||
|
||||
self.audio = numpy.empty(
|
||||
(self.chunks_to_record * self.BUFFERSIZE), dtype=numpy.int16)
|
||||
|
||||
def close(self):
|
||||
print("pyAudioRecorder closed")
|
||||
self.kill_threads = True
|
||||
self.p.close(self.in_stream)
|
||||
|
||||
### RECORDING AUDIO ###
|
||||
|
||||
def get_audio(self):
|
||||
"""get a single buffer size worth of audio."""
|
||||
audio_string = self.in_stream.read(self.BUFFERSIZE)
|
||||
return numpy.fromstring(audio_string, dtype=numpy.int16)
|
||||
|
||||
def record(self):
|
||||
while not self.kill_threads:
|
||||
for i in range(self.chunks_to_record):
|
||||
self.audio[i*self.BUFFERSIZE:(i+1)
|
||||
* self.BUFFERSIZE] = self.get_audio()
|
||||
self.has_new_audio = True
|
||||
|
||||
def start(self):
|
||||
print("pyAudioRecorder started")
|
||||
self.t = threading.Thread(target=self.record)
|
||||
self.t.start()
|
||||
|
||||
### MATH ###
|
||||
|
||||
def downsample(self, data, mult):
|
||||
"""Given 1D data, return the binned average."""
|
||||
overhang = len(data) % mult
|
||||
if overhang:
|
||||
data = data[:-overhang]
|
||||
data = numpy.reshape(data, (len(data) / mult, mult))
|
||||
data = numpy.average(data, 1)
|
||||
return data
|
||||
|
||||
def fft(self, data=None, trim_by=10, log_scale=False, div_by=100):
|
||||
if not data:
|
||||
data = self.audio.flatten()
|
||||
left, right = numpy.split(numpy.abs(numpy.fft.fft(data)), 2)
|
||||
ys = numpy.add(left, right[::-1])
|
||||
if log_scale:
|
||||
ys = numpy.multiply(20, numpy.log10(ys))
|
||||
xs = numpy.arange(self.BUFFERSIZE/2, dtype=float)
|
||||
if trim_by:
|
||||
i = int((self.BUFFERSIZE/2) / trim_by)
|
||||
ys = ys[:i]
|
||||
xs = xs[:i] * self.RATE / self.BUFFERSIZE
|
||||
if div_by:
|
||||
ys = ys / float(div_by)
|
||||
return xs, ys
|
||||
|
||||
### multithreading things ###
|
||||
|
||||
def registerRecorderClient(self,recorderClient):
|
||||
self.recorderClients.append(recorderClient)
|
||||
|
||||
def unregisterRecorderClient(self,recorderClient):
|
||||
self.recorderClients.remove(recorderClient)
|
||||
|
||||
class recorderClient(threading.Thread):
|
||||
def __init__():
|
||||
# toggle to true on beat, false when client got the value
|
||||
self.onBeat = False
|
||||
# when registering the client i want to be able to define how long the avg list should be
|
||||
self.y_max_freq_avg_list = []
|
||||
@@ -0,0 +1,54 @@
|
||||
from rgbUtils.RGBStrip import RGBStrip
|
||||
import time
|
||||
import threading
|
||||
import json
|
||||
class rgbStripController(threading.Thread):
|
||||
def __init__(self):
|
||||
threading.Thread.__init__(self)
|
||||
self.rgbStrips = []
|
||||
self.onRGBStripRegisteredHandler = []
|
||||
self.onRGBStripUnRegisteredHandler = []
|
||||
|
||||
def registerRGBStrip(self,rgbStripName,onValuesUpdateHandler):
|
||||
# maybe we can use an unique id if the strip reconnects later, eg push the uid
|
||||
# to the client on first connect and if he reconnects he sould send it back again.
|
||||
# the wmos could use the mac adress, if there is a python script it can save the uid
|
||||
# in a file or so.
|
||||
strip = RGBStrip(rgbStripName,onValuesUpdateHandler)
|
||||
self.rgbStrips.append(strip)
|
||||
self.noticeRGBStripRegisteredHandler(strip)
|
||||
return strip
|
||||
|
||||
def unregisterRGBStrip(self,strip):
|
||||
self.rgbStrips.remove(strip)
|
||||
self.noticeRGBStripUnRegisteredHandler(strip)
|
||||
|
||||
# returns all registered rgbStips
|
||||
def getRGBStrips(self):
|
||||
return self.rgbStrips
|
||||
|
||||
# inform all onRGBStripRegisteredHandler about the new RGBStrip
|
||||
def noticeRGBStripRegisteredHandler(self,rgbStrip):
|
||||
for hander in self.onRGBStripRegisteredHandler:
|
||||
hander(rgbStrip)
|
||||
|
||||
# add onRGBStripRegisteredHandler
|
||||
def addOnRGBStripRegisteredHandler(self, function):
|
||||
self.onRGBStripRegisteredHandler.append(function)
|
||||
|
||||
# remove onRGBStripRegisteredHandler
|
||||
def removeOnRGBStripRegisteredHandler(self, function):
|
||||
self.onRGBStripRegisteredHandler.remove(function)
|
||||
|
||||
# inform all onRGBStripUnRegisteredHandder about the removed RGBStrip
|
||||
def noticeRGBStripUnRegisteredHandler(self,rgbStrip):
|
||||
for hander in self.onRGBStripUnRegisteredHandler:
|
||||
hander(rgbStrip)
|
||||
|
||||
# add onRGBStripUnRegisteredHandler
|
||||
def addOnRGBStripUnRegisteredHandler(self, function):
|
||||
self.onRGBStripUnRegisteredHandler.append(function)
|
||||
|
||||
# remove onRGBStripUnRegisteredHandler
|
||||
def removeOnRGBStripUnRegisteredHandler(self, function):
|
||||
self.onRGBStripUnRegisteredHandler.remove(function)
|
||||
@@ -0,0 +1,14 @@
|
||||
"""
|
||||
Convert the rgbStripController function outputs to json format
|
||||
"""
|
||||
|
||||
# get the color values of a single led in json format
|
||||
def getRGBData(rgbStrip,led = 0):
|
||||
return {'led': led, 'red': rgbStrip.getData()[0][led], 'green': rgbStrip.getData()[1][led], 'blue': rgbStrip.getData()[2][led]}
|
||||
|
||||
# return json of all configured rgbStrips
|
||||
def getRGBStrips(rgbStripController):
|
||||
result = {}
|
||||
for rgbStrip in rgbStripController.getRGBStrips():
|
||||
result[rgbStrip.STRIP_UID] = {'index': rgbStrip.STRIP_UID, 'name': rgbStrip.STRIP_NAME}
|
||||
return result
|
||||
Reference in New Issue
Block a user