X-Git-Url: https://git.novaco.in/?p=electrum-server.git;a=blobdiff_plain;f=server.py;h=d41bcf35e4af8362455304eb527a67527b60ae5b;hp=666d5d446c23abc80de9225955b4c3f9c580cf08;hb=4ce69b7ea24ead59ebbcc7ed335ea9762ae3724b;hpb=aaf7eaa6b7e9cf1ef6e7afb4c2823912402f85f8 diff --git a/server.py b/server.py index 666d5d4..d41bcf3 100755 --- a/server.py +++ b/server.py @@ -15,572 +15,185 @@ # License along with this program. If not, see # . -""" -Todo: - * server should check and return bitcoind status.. - * improve txpoint sorting - * command to check cache - - mempool transactions do not need to be added to the database; it slows it down -""" - -import abe_backend - - - - -import time, json, socket, operator, thread, ast, sys, re, traceback import ConfigParser -from json import dumps, loads -import urllib - - -config = ConfigParser.ConfigParser() -# set some defaults, which will be overwritten by the config file -config.add_section('server') -config.set('server','banner', 'Welcome to Electrum!') -config.set('server', 'host', 'localhost') -config.set('server', 'port', '50000') -config.set('server', 'password', '') -config.set('server', 'irc', 'yes') -config.set('server', 'ircname', 'Electrum server') -config.add_section('database') -config.set('database', 'type', 'psycopg2') -config.set('database', 'database', 'abe') - -try: - f = open('/etc/electrum.conf','r') - config.readfp(f) - f.close() -except: - print "Could not read electrum.conf. I will use the default values." - -try: - f = open('/etc/electrum.banner','r') - config.set('server','banner', f.read()) - f.close() -except: - pass +import logging +import socket +import sys +import time +import threading +import traceback +import json -password = config.get('server','password') +logging.basicConfig() -stopping = False -block_number = -1 -sessions = {} -sessions_sub_numblocks = {} # sessions that have subscribed to the service +if sys.maxsize <= 2**32: + print "Warning: it looks like you are using a 32bit system. You may experience crashes caused by mmap" -m_sessions = [{}] # served by http - -peer_list = {} - -from Queue import Queue -input_queue = Queue() -output_queue = Queue() +def attempt_read_config(config, filename): + try: + with open(filename, 'r') as f: + config.readfp(f) + except IOError: + pass + + +def create_config(): + config = ConfigParser.ConfigParser() + # set some defaults, which will be overwritten by the config file + config.add_section('server') + config.set('server', 'banner', 'Welcome to Electrum!') + config.set('server', 'host', 'localhost') + config.set('server', 'report_host', '') + config.set('server', 'stratum_tcp_port', '50001') + config.set('server', 'stratum_http_port', '8081') + config.set('server', 'stratum_tcp_ssl_port', '') + config.set('server', 'stratum_http_ssl_port', '') + config.set('server', 'report_stratum_tcp_port', '') + config.set('server', 'report_stratum_http_port', '') + config.set('server', 'report_stratum_tcp_ssl_port', '') + config.set('server', 'report_stratum_http_ssl_port', '') + config.set('server', 'ssl_certfile', '') + config.set('server', 'ssl_keyfile', '') + config.set('server', 'password', '') + config.set('server', 'irc', 'yes') + config.set('server', 'irc_nick', '') + config.set('server', 'coin', '') + config.set('server', 'datadir', '') + + # use leveldb as default + config.set('server', 'backend', 'leveldb') + config.add_section('leveldb') + config.set('leveldb', 'path', '/dev/shm/electrum_db') + config.set('leveldb', 'pruning_limit', '100') + + for path in ('/etc/', ''): + filename = path + 'electrum.conf' + attempt_read_config(config, filename) + try: + with open('/etc/electrum.banner', 'r') as f: + config.set('server', 'banner', f.read()) + except IOError: + pass + return config -def random_string(N): - import random, string - return ''.join(random.choice(string.ascii_uppercase + string.digits) for x in range(N)) +def run_rpc_command(command, stratum_tcp_port): + try: + s = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + s.connect((host, int(stratum_tcp_port))) + except: + print "cannot connect to server." + return + + method = 'server.' + command + params = [password] + if len(sys.argv) > 2: + params.append(sys.argv[2]) - -def cmd_stop(_,__,pw): - global stopping - if password == pw: - stopping = True - return 'ok' - else: - return 'wrong password' - -def cmd_load(_,__,pw): - if password == pw: - return repr( len(sessions) ) - else: - return 'wrong password' - - - - - -def modified_addresses(a_session): - #t1 = time.time() - import copy - session = copy.deepcopy(a_session) - addresses = session['addresses'] - session['last_time'] = time.time() - ret = {} - k = 0 - for addr in addresses: - status = store.get_status( addr ) - msg_id, last_status = addresses.get( addr ) - if last_status != status: - addresses[addr] = msg_id, status - ret[addr] = status - - #t2 = time.time() - t1 - #if t2 > 10: print "high load:", session_id, "%d/%d"%(k,len(addresses)), t2 - return ret, addresses - - -def poll_session(session_id): - # native - session = sessions.get(session_id) - if session is None: - print time.asctime(), "session not found", session_id - return -1, {} - else: - sessions[session_id]['last_time'] = time.time() - ret, addresses = modified_addresses(session) - if ret: sessions[session_id]['addresses'] = addresses - return repr( (block_number,ret)) - - -def poll_session_json(session_id, message_id): - session = m_sessions[0].get(session_id) - if session is None: - raise BaseException("session not found %s"%session_id) + request = json.dumps({'id': 0, 'method': method, 'params': params}) + s.send(request + '\n') + msg = '' + while True: + o = s.recv(1024) + if not o: break + msg += o + if msg.find('\n') != -1: + break + s.close() + r = json.loads(msg).get('result') + + if command == 'info': + now = time.time() + print 'type address sub version time' + for item in r: + print '%4s %21s %3s %7s %.2f' % (item.get('name'), + item.get('address'), + item.get('subscriptions'), + item.get('version'), + (now - item.get('time')), + ) else: - m_sessions[0][session_id]['last_time'] = time.time() - out = [] - ret, addresses = modified_addresses(session) - if ret: - m_sessions[0][session_id]['addresses'] = addresses - for addr in ret: - msg_id, status = addresses[addr] - out.append( { 'id':msg_id, 'result':status } ) - - msg_id, last_nb = session.get('numblocks') - if last_nb: - if last_nb != block_number: - m_sessions[0][session_id]['numblocks'] = msg_id, block_number - out.append( {'id':msg_id, 'result':block_number} ) - - return out - - -def do_update_address(addr): - # an address was involved in a transaction; we check if it was subscribed to in a session - # the address can be subscribed in several sessions; the cache should ensure that we don't do redundant requests - - for session_id in sessions.keys(): - session = sessions[session_id] - if session.get('type') != 'persistent': continue - addresses = session['addresses'].keys() - - if addr in addresses: - status = store.get_status( addr ) - message_id, last_status = session['addresses'][addr] - if last_status != status: - #print "sending new status for %s:"%addr, status - send_status(session_id,message_id,addr,status) - sessions[session_id]['addresses'][addr] = (message_id,status) - - - -def send_numblocks(session_id): - message_id = sessions_sub_numblocks[session_id] - out = json.dumps( {'id':message_id, 'result':block_number} ) - output_queue.put((session_id, out)) - -def send_status(session_id, message_id, address, status): - out = json.dumps( { 'id':message_id, 'result':status } ) - output_queue.put((session_id, out)) - -def address_get_history_json(_,message_id,address): - return store.get_history(address) - -def subscribe_to_numblocks(session_id, message_id): - sessions_sub_numblocks[session_id] = message_id - send_numblocks(session_id) - -def subscribe_to_numblocks_json(session_id, message_id): - global m_sessions - m_sessions[0][session_id]['numblocks'] = message_id,block_number - return block_number - -def subscribe_to_address(session_id, message_id, address): - status = store.get_status(address) - sessions[session_id]['addresses'][address] = (message_id, status) - sessions[session_id]['last_time'] = time.time() - send_status(session_id, message_id, address, status) - -def add_address_to_session_json(session_id, message_id, address): - global m_sessions - sessions = m_sessions[0] - status = store.get_status(address) - sessions[session_id]['addresses'][address] = (message_id, status) - sessions[session_id]['last_time'] = time.time() - m_sessions[0] = sessions - return status - -def add_address_to_session(session_id, address): - status = store.get_status(address) - sessions[session_id]['addresses'][address] = ("", status) - sessions[session_id]['last_time'] = time.time() - return status - -def new_session(version, addresses): - session_id = random_string(10) - sessions[session_id] = { 'addresses':{}, 'version':version } - for a in addresses: - sessions[session_id]['addresses'][a] = ('','') - out = repr( (session_id, config.get('server','banner').replace('\\n','\n') ) ) - sessions[session_id]['last_time'] = time.time() - return out - - -def client_version_json(session_id, _, version): - global m_sessions - sessions = m_sessions[0] - sessions[session_id]['version'] = version - m_sessions[0] = sessions - -def create_session_json(_, __): - sessions = m_sessions[0] - session_id = random_string(10) - print "creating session", session_id - sessions[session_id] = { 'addresses':{}, 'numblocks':('','') } - sessions[session_id]['last_time'] = time.time() - m_sessions[0] = sessions - return session_id - - - -def get_banner(_,__): - return config.get('server','banner').replace('\\n','\n') - -def update_session(session_id,addresses): - """deprecated in 0.42""" - sessions[session_id]['addresses'] = {} - for a in addresses: - sessions[session_id]['addresses'][a] = '' - sessions[session_id]['last_time'] = time.time() - return 'ok' - -def native_server_thread(): - s = socket.socket(socket.AF_INET, socket.SOCK_STREAM) - s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) - s.bind((config.get('server','host'), config.getint('server','port'))) - s.listen(1) - while not stopping: - conn, addr = s.accept() - try: - thread.start_new_thread(native_client_thread, (addr, conn,)) - except: - # can't start new thread if there is no memory.. - traceback.print_exc(file=sys.stdout) - - -def native_client_thread(ipaddr,conn): - #print "client thread", ipaddr - try: - ipaddr = ipaddr[0] - msg = '' - while 1: - d = conn.recv(1024) - msg += d - if not d: - break - if '#' in msg: - msg = msg.split('#', 1)[0] - break - try: - cmd, data = ast.literal_eval(msg) - except: - print "syntax error", repr(msg), ipaddr - conn.close() - return - - out = do_command(cmd, data, ipaddr) - if out: - #print ipaddr, cmd, len(out) - try: - conn.send(out) - except: - print "error, could not send" - - finally: - conn.close() - + print r -def timestr(): - return time.strftime("[%d/%m/%Y-%H:%M:%S]") -# used by the native handler -def do_command(cmd, data, ipaddr): - - if cmd=='b': - out = "%d"%block_number - - elif cmd in ['session','new_session']: - try: - if cmd == 'session': - addresses = ast.literal_eval(data) - version = "old" - else: - version, addresses = ast.literal_eval(data) - if version[0]=="0": version = "v" + version - except: - print "error", data - return None - print timestr(), "new session", ipaddr, addresses[0] if addresses else addresses, len(addresses), version - out = new_session(version, addresses) - - elif cmd=='address.subscribe': - try: - session_id, addr = ast.literal_eval(data) - except: - traceback.print_exc(file=sys.stdout) - print data - return None - out = add_address_to_session(session_id,addr) - - elif cmd=='update_session': - try: - session_id, addresses = ast.literal_eval(data) - except: - traceback.print_exc(file=sys.stdout) - return None - print timestr(), "update session", ipaddr, addresses[0] if addresses else addresses, len(addresses) - out = update_session(session_id,addresses) - - elif cmd=='poll': - out = poll_session(data) - - elif cmd == 'h': - # history - address = data - out = repr( store.get_history( address ) ) - - elif cmd == 'load': - out = cmd_load(None,None,data) - - elif cmd =='tx': - out = store.send_tx(data) - print timestr(), "sent tx:", ipaddr, out - - elif cmd == 'stop': - out = cmd_stop(data) +if __name__ == '__main__': + config = create_config() + password = config.get('server', 'password') + host = config.get('server', 'host') + stratum_tcp_port = config.get('server', 'stratum_tcp_port') + stratum_http_port = config.get('server', 'stratum_http_port') + stratum_tcp_ssl_port = config.get('server', 'stratum_tcp_ssl_port') + stratum_http_ssl_port = config.get('server', 'stratum_http_ssl_port') + ssl_certfile = config.get('server', 'ssl_certfile') + ssl_keyfile = config.get('server', 'ssl_keyfile') + + if stratum_tcp_ssl_port or stratum_http_ssl_port: + assert ssl_certfile and ssl_keyfile + + if len(sys.argv) > 1: + run_rpc_command(sys.argv[1], stratum_tcp_port) + sys.exit(0) - elif cmd == 'peers': - out = repr(peer_list.values()) + from processor import Dispatcher, print_log + from backends.irc import ServerProcessor + backend_name = config.get('server', 'backend') + if backend_name == 'libbitcoin': + from backends.libbitcoin import BlockchainProcessor + elif backend_name == 'leveldb': + from backends.bitcoind import BlockchainProcessor else: - out = None - - return out - - -def clean_session_thread(): - while not stopping: - time.sleep(30) - t = time.time() - for k,s in sessions.items(): - if s.get('type') == 'persistent': continue - t0 = s['last_time'] - if t - t0 > 5*60: - sessions.pop(k) - print "lost session", k - - -#################################################################### + print "Unknown backend '%s' specified\n" % backend_name + sys.exit(1) + for i in xrange(5): + print "" + print_log("Starting Electrum server on", host) -import stratum + # Create hub + dispatcher = Dispatcher(config) + shared = dispatcher.shared -class AbeProcessor(stratum.Processor): - def process(self,session,request): - message_id = request['id'] - method = request['method'] - params = request.get('params',[]) + # Create and register processors + chain_proc = BlockchainProcessor(config, shared) + dispatcher.register('blockchain', chain_proc) - result = '' - if method == 'numblocks.subscribe': - session.subscribe_to_numblocks(message_id) - result = block_number - elif method == 'address.subscribe': - address = params[0] - status = store.get_status(address) - session.subscribe_to_address(address,message_id,status) - result = status - elif method == 'client.version': - session.version = params[0] - elif method == 'server.banner': - result = config.get('server','banner').replace('\\n','\n') - elif method == 'server.peers': - result = peer_list.values() - elif method == 'address.get_history': - address = params[0] - result = store.get_history( address ) - elif method == 'transaction.broadcast': - txo = store.send_tx(params[0]) - print "sent tx:", txo - result = txo - else: - print "unknown method", request + server_proc = ServerProcessor(config) + dispatcher.register('server', server_proc) - if result!='': - response = { 'id':message_id, 'result':result } - self.push_response(session,response) - - - -#################################################################### - - - - - -def irc_thread(): - global peer_list - NICK = 'E_'+random_string(10) - while not stopping: + transports = [] + # Create various transports we need + if stratum_tcp_port: + from transports.stratum_tcp import TcpServer + tcp_server = TcpServer(dispatcher, host, int(stratum_tcp_port), False, None, None) + transports.append(tcp_server) + + if stratum_tcp_ssl_port: + from transports.stratum_tcp import TcpServer + tcp_server = TcpServer(dispatcher, host, int(stratum_tcp_ssl_port), True, ssl_certfile, ssl_keyfile) + transports.append(tcp_server) + + if stratum_http_port: + from transports.stratum_http import HttpServer + http_server = HttpServer(dispatcher, host, int(stratum_http_port), False, None, None) + transports.append(http_server) + + if stratum_http_ssl_port: + from transports.stratum_http import HttpServer + http_server = HttpServer(dispatcher, host, int(stratum_http_ssl_port), True, ssl_certfile, ssl_keyfile) + transports.append(http_server) + + for server in transports: + server.start() + + while not shared.stopped(): try: - s = socket.socket() - s.connect(('irc.freenode.net', 6667)) - s.send('USER electrum 0 * :'+config.get('server','host')+' '+config.get('server','ircname')+'\n') - s.send('NICK '+NICK+'\n') - s.send('JOIN #electrum\n') - sf = s.makefile('r', 0) - t = 0 - while not stopping: - line = sf.readline() - line = line.rstrip('\r\n') - line = line.split() - if line[0]=='PING': - s.send('PONG '+line[1]+'\n') - elif '353' in line: # answer to /names - k = line.index('353') - for item in line[k+1:]: - if item[0:2] == 'E_': - s.send('WHO %s\n'%item) - elif '352' in line: # answer to /who - # warning: this is a horrible hack which apparently works - k = line.index('352') - ip = line[k+4] - ip = socket.gethostbyname(ip) - name = line[k+6] - host = line[k+9] - peer_list[name] = (ip,host) - if time.time() - t > 5*60: - s.send('NAMES #electrum\n') - t = time.time() - peer_list = {} + time.sleep(1) except: - traceback.print_exc(file=sys.stdout) - finally: - sf.close() - s.close() - - -def get_peers_json(_,__): - return peer_list.values() - -def http_server_thread(): - # see http://code.google.com/p/jsonrpclib/ - from SocketServer import ThreadingMixIn - from StratumJSONRPCServer import StratumJSONRPCServer - class StratumThreadedJSONRPCServer(ThreadingMixIn, StratumJSONRPCServer): pass - server = StratumThreadedJSONRPCServer(( config.get('server','host'), 8081)) - server.register_function(get_peers_json, 'server.peers') - server.register_function(cmd_stop, 'stop') - server.register_function(cmd_load, 'load') - server.register_function(get_banner, 'server.banner') - server.register_function(lambda a,b,c: store.send_tx(c), 'transaction.broadcast') - server.register_function(address_get_history_json, 'address.get_history') - server.register_function(add_address_to_session_json, 'address.subscribe') - server.register_function(subscribe_to_numblocks_json, 'numblocks.subscribe') - server.register_function(client_version_json, 'client.version') - server.register_function(create_session_json, 'session.create') # internal message (not part of protocol) - server.register_function(poll_session_json, 'session.poll') # internal message (not part of protocol) - server.serve_forever() - - -if __name__ == '__main__': - - if len(sys.argv)>1: - import jsonrpclib - server = jsonrpclib.Server('http://%s:8081'%config.get('server','host')) - cmd = sys.argv[1] - if cmd == 'load': - out = server.load(password) - elif cmd == 'peers': - out = server.server.peers() - elif cmd == 'stop': - out = server.stop(password) - elif cmd == 'clear_cache': - out = server.clear_cache(password) - elif cmd == 'get_cache': - out = server.get_cache(password,sys.argv[2]) - elif cmd == 'h': - out = server.address.get_history(sys.argv[2]) - elif cmd == 'tx': - out = server.transaction.broadcast(sys.argv[2]) - elif cmd == 'b': - out = server.numblocks.subscribe() - else: - out = "Unknown command: '%s'" % cmd - print out - sys.exit(0) - - # backend - store = abe_backend.AbeStore(config) - - # supported protocols - thread.start_new_thread(native_server_thread, ()) - - thread.start_new_thread(http_server_thread, ()) - thread.start_new_thread(clean_session_thread, ()) - - #tcp stratum - stratum_processor = AbeProcessor() - shared = stratum.Shared() - # Bind shared to processor since constructor is user defined - stratum_processor.shared = shared - stratum_processor.start() - # Create various transports we need - server = stratum.TcpServer(shared, stratum_processor, "ecdsa.org",50001) - server.start() - - if (config.get('server','irc') == 'yes' ): - thread.start_new_thread(irc_thread, ()) - - print "starting Electrum server" - - old_block_number = None - while not stopping: - block_number = store.main_iteration() - - if block_number != old_block_number: - old_block_number = block_number - for session_id in sessions_sub_numblocks.keys(): - send_numblocks(session_id) - - for session in stratum_processor.sessions: - if session.numblocks_sub is not None: - response = { 'id':session.numblocks_sub, 'result':block_number } - stratum_processor.push_response(session,response) - - while True: - try: - addr = store.address_queue.get(False) - except: - break - do_update_address(addr) - - for session in stratum_processor.sessions: - m = session.addresses_sub.get(addr) - if m: - status = store.get_status( addr ) - message_id, last_status = m - if status != last_status: - session.subscribe_to_address(message_id, status) - response = { 'id':message_id, 'result':status } - stratum_processor.push_response(session,response) - - time.sleep(10) - print "server stopped" + shared.stop() + print_log("Electrum Server stopped")